# 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

1. **Tạo một dự án Spring Boot** Sử dụng Spring Initializr ([https://start.spring.io/](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`:
    
    ```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`](http://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`](http://application.properties):
    
    ```java
    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`:
    
    ```java
    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`:
    
    ```java
    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`:
    
    ```java
    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:
    
    ```sh
    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`:
    
    ```sh
    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`:
        
        ```sh
        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`:
        
        ```java
        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`](http://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.
