Bài 14: Stream Processing
Mục tiêu
Sử dụng Kafka Streams API với Spring Boot.
Nội dung chính
Giới thiệu Kafka Streams API
Kafka Streams API là một thư viện mạnh mẽ và linh hoạt để xử lý dữ liệu streaming trong Kafka. Nó cung cấp các khả năng để xử lý dữ liệu theo thời gian thực, từ việc lọc, tổng hợp, chuyển đổi dữ liệu đến việc phân tích dữ liệu.
Các khái niệm cơ bản và cấu trúc của Kafka Streams
Một số khái niệm cơ bản trong Kafka Streams bao gồm:
Stream: Một dòng liên tục của các bản ghi.
Topology: Một đồ thị định nghĩa cách xử lý dữ liệu streaming.
KStream: Một abstract stream của các bản ghi từ Kafka.
KTable: Một abstract table của các bản ghi từ Kafka.
Processor API: API để định nghĩa các thao tác xử lý tùy chỉnh.
Tích hợp Kafka Streams với Spring Boot
Spring Boot cung cấp tích hợp với Kafka Streams thông qua các dependency và cấu hình dễ dàng. Bạn có thể sử dụng Kafka Streams API trong các ứng dụng Spring Boot để thực hiện các thao tác xử lý dữ liệu streaming.
Thực hành
Viết ứng dụng Spring Boot sử dụng Kafka Streams API
Tạo một dự án Spring Boot Sử dụng Spring Initializr (https://start.spring.io/) để tạo một dự án Spring Boot mới. Chọn các thông tin cơ bản như
Group,Artifact, và thêm các dependency cần thiết nhưSpring Web,Spring Kafka,Spring Kafka Streams.Cấu hình
pom.xmlThêm các dependency cần thiết vào tệppom.xml:<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka-streams</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> </dependencies>Cấu hình
application.propertiesCấu hình các thông số kết nối đến Kafka broker trong tệpsrc/main/resources/application.properties:spring.kafka.bootstrap-servers=localhost:9092 spring.kafka.streams.application-id=my-streams-app spring.kafka.streams.bootstrap-servers=localhost:9092Tạo Kafka Streams Configuration Tạo lớp cấu hình
KafkaStreamsConfigtrong thư mụcsrc/main/java/com/example/config:import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.kstream.KStream; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.annotation.EnableKafkaStreams; import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.config.StreamsBuilderFactoryBean; import java.util.HashMap; import java.util.Map; @Configuration @EnableKafkaStreams public class KafkaStreamsConfig { @Bean(name = KafkaStreamsDefaultConfiguration.DEFAULT_STREAMS_CONFIG_BEAN_NAME) public KafkaStreamsConfiguration kStreamsConfigs() { Map<String, Object> props = new HashMap<>(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "my-streams-app"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); return new KafkaStreamsConfiguration(props); } @Bean public KStream<String, String> kStream(StreamsBuilder streamsBuilder) { KStream<String, String> stream = streamsBuilder.stream("input-topic"); stream.filter((key, value) -> value.contains("filter")) .mapValues(value -> value.toUpperCase()) .to("output-topic"); return stream; } }Tạo Kafka Producer Service để gửi thông điệp Tạo một lớp service để gửi thông điệp tới Kafka trong thư mục
src/main/java/com/example/service:import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; @Service public class KafkaProducerService { private final KafkaTemplate<String, String> kafkaTemplate; public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public void sendMessage(String topic, String message) { kafkaTemplate.send(topic, message); } }Tạo REST Controller để gửi thông điệp Tạo một REST controller để gửi thông điệp qua API trong thư mục
src/main/java/com/example/controller:import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; @RestController public class KafkaController { private final KafkaProducerService kafkaProducerService; public KafkaController(KafkaProducerService kafkaProducerService) { this.kafkaProducerService = kafkaProducerService; } @PostMapping("/send") public String sendMessage(@RequestParam String topic, @RequestParam String message) { kafkaProducerService.sendMessage(topic, message); return "Message sent!"; } }Chạy ứng dụng Spring Boot Chạy ứng dụng Spring Boot bằng cách sử dụng IDE hoặc dòng lệnh:
mvn spring-boot:runThực hiện các thao tác stream processing cơ bản Sử dụng công cụ như Postman hoặc curl để gửi request đến endpoint
/send:curl -X POST "http://localhost:8080/send?topic=input-topic&message=filterThisMessage" curl -X POST "http://localhost:8080/send?topic=input-topic&message=ignoreThisMessage"Kiểm tra đầu ra của console và sử dụng
kafka-console-consumerđể nhận thông điệp từoutput-topic:bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic output-topic --from-beginningBạn sẽ thấy thông điệp "FILTERTHISMESSAGE" được gửi tới
output-topic:FILTERTHISMESSAGE
Câu hỏi củng cố kiến thức
Câu hỏi: Kafka Streams API là gì? Trả lời: Kafka Streams API là một thư viện mạnh mẽ và linh hoạt để xử lý dữ liệu streaming trong Kafka. Nó cung cấp các khả năng để xử lý dữ liệu theo thời gian thực, từ việc lọc, tổng hợp, chuyển đổi dữ liệu đến việc phân tích dữ liệu.
Câu hỏi: Các khái niệm cơ bản trong Kafka Streams là gì? Trả lời: Các khái niệm cơ bản trong Kafka Streams bao gồm Stream, Topology, KStream, KTable, và Processor API.
Câu hỏi: Làm thế nào để tích hợp Kafka Streams với Spring Boot? Trả lời: Để tích hợp Kafka Streams với Spring Boot, bạn cần thêm dependency
spring-kafka-streamsvào dự án và cấu hình các thông số cần thiết trongapplication.properties. Sau đó, tạo các bean cấu hình cho Kafka Streams và định nghĩa các thao tác stream processing trong các lớp cấu hình.Câu hỏi: Làm thế nào để thực hiện các thao tác stream processing cơ bản trong Spring Boot? Trả lời: Bạn có thể sử dụng
StreamsBuilderđể định nghĩa các thao tác stream processing cơ bản như filter, map, và to. Sau đó, sử dụngKStreamđể thực hiện các thao tác này trên các stream từ Kafka.