Skip to main content

Command Palette

Search for a command to run...

Bài 6: Consumer API

Published
4 min readView as Markdown

Mục tiêu

Hiểu và sử dụng Consumer API để nhận thông điệp từ Kafka.

Nội dung chính

Khái niệm về Consumer trong Kafka

Consumer là các ứng dụng hoặc dịch vụ nhận và xử lý các message từ các topic trong Kafka. Một consumer có thể đăng ký nhận message từ một hoặc nhiều topic và Kafka sẽ tự động gửi các message từ các partition của topic đó tới consumer.

Consumer Group: Các consumer thường hoạt động theo nhóm gọi là consumer group. Mỗi message trong một partition sẽ được xử lý bởi chỉ một consumer trong nhóm đó. Điều này đảm bảo rằng mỗi message chỉ được xử lý một lần trong một nhóm, nhưng có thể được xử lý bởi nhiều consumer nếu có nhiều nhóm.

Các cấu hình cơ bản của Consumer

Khi cấu hình một Kafka consumer, có một số thuộc tính quan trọng cần thiết:

  1. bootstrap.servers: Địa chỉ của một hoặc nhiều Kafka broker (ví dụ: localhost:9092).

  2. group.id: ID của consumer group mà consumer này thuộc về.

  3. key.deserializer: Lớp dùng để deserialize khóa của message (ví dụ: org.apache.kafka.common.serialization.StringDeserializer).

  4. value.deserializer: Lớp dùng để deserialize giá trị của message (ví dụ: org.apache.kafka.common.serialization.StringDeserializer).

  5. auto.offset.reset: Xác định cách xử lý khi không tìm thấy offset trước đó cho một partition (ví dụ: earliest, latest).

Nhận thông điệp từ Kafka đồng bộ và không đồng bộ

  • Nhận đồng bộ: Consumer chủ động gọi để nhận message và xử lý từng message một cách tuần tự.

  • Nhận không đồng bộ: Consumer đăng ký một callback để xử lý message khi chúng đến mà không cần chờ đợi tuần tự.

Thực hành

Viết một ứng dụng Java sử dụng Consumer API để nhận thông điệp từ Kafka

  1. Tạo Maven Project

     <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
         <modelVersion>4.0.0</modelVersion>
         <groupId>com.example</groupId>
         <artifactId>kafka-consumer-example</artifactId>
         <version>1.0-SNAPSHOT</version>
         <dependencies>
             <dependency>
                 <groupId>org.apache.kafka</groupId>
                 <artifactId>kafka-clients</artifactId>
                 <version>2.8.0</version>
             </dependency>
         </dependencies>
     </project>
    
  2. Cấu hình Kafka Consumer

     import org.apache.kafka.clients.consumer.ConsumerConfig;
     import org.apache.kafka.clients.consumer.ConsumerRecord;
     import org.apache.kafka.clients.consumer.ConsumerRecords;
     import org.apache.kafka.clients.consumer.KafkaConsumer;
     import org.apache.kafka.common.serialization.StringDeserializer;
    
     import java.time.Duration;
     import java.util.Collections;
     import java.util.Properties;
    
     public class KafkaConsumerExample {
         public static void main(String[] args) {
             Properties properties = new Properties();
             properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
             properties.put(ConsumerConfig.GROUP_ID_CONFIG, "my-group");
             properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
             properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
             properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    
             KafkaConsumer<String, String> consumer = new KafkaConsumer<>(properties);
             consumer.subscribe(Collections.singletonList("my-topic"));
    
             try {
                 while (true) {
                     ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                     for (ConsumerRecord<String, String> record : records) {
                         System.out.printf("Consumed message: topic = %s, partition = %d, offset = %d, key = %s, value = %s%n",
                                 record.topic(), record.partition(), record.offset(), record.key(), record.value());
                     }
                 }
             } finally {
                 consumer.close();
             }
         }
     }
    

Xác minh thông điệp đã được nhận thành công

  1. Chạy Consumer:

    • Chạy chương trình Java trên để bắt đầu nhận message từ Kafka.
  2. Gửi Message bằng Kafka Console Producer:

    • Mở một terminal mới và sử dụng lệnh sau để gửi message vào topic:

        bin/kafka-console-producer.sh --broker-list localhost:9092 --topic my-topic
      
    • Nhập một vài message và nhấn Enter.

  3. Xác minh trên Consumer:

    • Quan sát đầu ra của chương trình consumer để đảm bảo rằng các message đã được nhận và hiển thị.

Giải thích chi tiết từng phần của mã Java

  1. Properties Configuration:

    • ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG: Địa chỉ của Kafka broker mà consumer sẽ kết nối.

    • ConsumerConfig.GROUP_ID_CONFIG: ID của consumer group mà consumer này sẽ tham gia.

    • ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG: Lớp dùng để deserialize khóa của message.

    • ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG: Lớp dùng để deserialize giá trị của message.

    • ConsumerConfig.AUTO_OFFSET_RESET_CONFIG: Xác định hành vi của consumer khi không tìm thấy offset trước đó (ví dụ: đọc từ đầu earliest).

  2. KafkaConsumer Initialization:

    • KafkaConsumer<String, String> consumer = new KafkaConsumer<>(properties);: Tạo một instance KafkaConsumer với cấu hình đã thiết lập.
  3. Subscription:

    • consumer.subscribe(Collections.singletonList("my-topic"));: Đăng ký consumer để nhận message từ topic "my-topic".
  4. Polling Loop:

    • ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));: Consumer sẽ lấy các message mới từ Kafka mỗi 100 milliseconds.

    • for (ConsumerRecord<String, String> record : records) { ... }: Lặp qua tất cả các record được lấy và xử lý chúng.

  5. Message Processing:

    • System.out.printf("Consumed message: topic = %s, partition = %d, offset = %d, key = %s, value = %s%n", record.topic(), record.partition(), record.offset(), record.key(), record.value());: Hiển thị thông tin chi tiết của mỗi message được nhận.

Kết luận

Với bài viết này, bạn đã hiểu cách cấu hình và sử dụng Consumer API trong Kafka để nhận message từ các topic. Bạn cũng đã thực hiện các thao tác nhận message đồng bộ và xác minh rằng message đã được nhận thành công. Trong bài tiếp theo, chúng ta sẽ đi sâu vào các kỹ thuật nâng cao hơn như xử lý message không đồng bộ và quản lý offset.

More from this blog

hoangkim

366 posts