Skip to main content

Command Palette

Search for a command to run...

Bài 14: Stream Processing

Published
4 min readView as Markdown

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

  1. 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.

  2. Cấu hình pom.xml Thêm các dependency cần thiết vào tệp pom.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>
    
  3. Cấu hình application.properties Cấu hình các thông số kết nối đến Kafka broker trong tệp src/main/resources/application.properties:

     spring.kafka.bootstrap-servers=localhost:9092
     spring.kafka.streams.application-id=my-streams-app
     spring.kafka.streams.bootstrap-servers=localhost:9092
    
  4. Tạo Kafka Streams Configuration Tạo lớp cấu hình KafkaStreamsConfig trong thư mục src/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;
         }
     }
    
  5. 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);
         }
     }
    
  6. 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!";
         }
     }
    
  7. 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:run
    
  8. Thự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-beginning
      
    • Bạ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

  1. 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.

  2. 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.

  3. 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-streams vào dự án và cấu hình các thông số cần thiết trong application.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.

  4. 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ụng KStream để thực hiện các thao tác này trên các stream từ Kafka.

More from this blog

hoangkim

366 posts