Phân tích và Triển khai Cơ chế Tiêu thụ Thủ công (Manual Consumer) trong Apache Kafka

Trong hệ sinh thái Apache Kafka, bên cạnh cơ chế Consumer Group tự động (High-Level API), hệ thống còn cung cấp khả năng kiểm soát chi tiết ở tầng thấp hơn đối với quá trình tiêu thụ tin nhắn. Trong các phiên bản Kafka cũ, tính năng này được gọi là Simple Consumer API. Ở các phiên bản hiện đại, nó được thay thế bằng KafkaConsumer với cơ chế gán partition thủ công (Manual Partition Assignment). Bài viết này sẽ phân tích các trường hợp sử dụng, thách thức và cách triển khai cơ chế tiêu thụ cấp thấp này.

1. Trường hợp sử dụng và Thách thức

Việc sử dụng API tiêu thụ cấp thấp mang lại quyền kiểm soát tối đa, phù hợp cho các kịch bản đặc thù:

  • Đọc lại tin nhắn (Replay): Cho phép đọc lại một thông điệp nhiều lần hoặc tua về quá khứ để xử lý lại dữ liệu.
  • Tiêu thụ một phần Partition: Chỉ định rõ ràng việc chỉ tiêu thụ một phân vùng (partition) cụ thể thay vì để Kafka tự động cân bằng.
  • Kiểm soát Offset và Giao dịch: Tự tay quản lý offset để tích hợp với các hệ thống giao dịch bên ngoài, đảm bảo ngữ nghĩa "exact-once" (xử lý chính xác một lần).

Tuy nhiên, quyền lực đi kèm với trách nhiệm. Khi không sử dụng Consumer Group, lập trình viên phải tự xử lý các thách thức sau:

  • Tự theo dõi và lưu trữ giá trị Offset.
  • Xử lý các lỗi kết nối mạng ở tầng thấp và tự xây dựng cơ chế chịu lỗi (Fault Tolerance).
  • Quản lý vòng đời của Consumer một cách chặt chẽ để tránh rò rỉ tài nguyên.

2. Cấu hình Dependency

Để triển khai, bạn cần thêm thư viện Kafka Client vào dự án. Đối với Maven, hãy thêm dependency sau vào pom.xml:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.6.0</version>
</dependency>

3. Triển khai Mã nguồn

Dưới đây là đoạn mã minh họa cách sử dụng KafkaConsumer để gán thủ công partition và quản lý offset. Cấu trúc mã đã được tái cấu trúc để tách biệt logic khởi tạo, vòng lặp tiêu thụ và xử lý ngoại lệ, giúp mã nguồn dễ bảo trì và mở rộng hơn.

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.TopicPartition;
import org.apache.kafka.common.errors.WakeupException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
import java.util.concurrent.atomic.AtomicBoolean;

public class ManualPartitionProcessor {

    private static final Logger logger = LoggerFactory.getLogger(ManualPartitionProcessor.class);
    private final AtomicBoolean isRunning = new AtomicBoolean(true);
    private KafkaConsumer<String, String> consumer;

    public void initialize(String bootstrapServers, String groupId) {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        // Tắt tự động commit để quản lý offset thủ công
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");

        this.consumer = new KafkaConsumer<>(props);
    }

    public void processPartition(String topic, int partition, long maxRecordsToProcess) {
        if (consumer == null) {
            throw new IllegalStateException("Consumer chưa được khởi tạo.");
        }

        TopicPartition targetPartition = new TopicPartition(topic, partition);
        // Gán thủ công partition, không tham gia Consumer Group
        consumer.assign(Collections.singletonList(targetPartition));
        
        // Tua về đầu partition để đọc lại toàn bộ (hoặc dùng seek() để đến offset cụ thể)
        consumer.seekToBeginning(Collections.singletonList(targetPartition));

        long processedCount = 0;
        try {
            while (isRunning.get() && processedCount < maxRecordsToProcess) {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
                
                for (ConsumerRecord<String, String> record : records) {
                    handleRecord(record);
                    processedCount++;
                    
                    if (processedCount >= maxRecordsToProcess) {
                        break;
                    }
                }
                
                // Commit offset thủ công sau mỗi đợt poll
                if (!records.isEmpty()) {
                    consumer.commitSync();
                }
            }
        } catch (WakeupException e) {
            logger.info("Nhận tín hiệu wake-up, đang đóng gracefully...");
        } finally {
            shutdown();
        }
    }

    private void handleRecord(ConsumerRecord<String, String> record) {
        logger.info("Đang xử lý tin nhắn - Topic: {}, Partition: {}, Offset: {}, Key: {}, Value: {}",
                record.topic(), record.partition(), record.offset(), record.key(), record.value());
        // Logic xử lý nghiệp vụ tại đây
    }

    public void shutdown() {
        isRunning.set(false);
        if (consumer != null) {
            consumer.wakeup();
            consumer.close();
            logger.info("Đã đóng kết nối Consumer.");
        }
    }

    public static void main(String[] args) {
        ManualPartitionProcessor processor = new ManualPartitionProcessor();
        processor.initialize("localhost:9092", "manual-processor-group");
        
        String targetTopic = "events-stream";
        int targetPartition = 0;
        long maxRecords = 1000;

        Runtime.getRuntime().addShutdownHook(new Thread(processor::shutdown));
        
        processor.processPartition(targetTopic, targetPartition, maxRecords);
    }
}

4. Phân tích Logic Mã nguồn

  • Tắt Auto-Commit: Thuộc tính ENABLE_AUTO_COMMIT_CONFIG được đặt thành false. Điều này ngăn Kafka tự động lưu offset, trao toàn quyền kiểm soát cho lập trình viên.
  • Manual Assignment: Phương thức assign() được sử dụng thay vì subscribe(). Điều này bypass hoàn toàn cơ chế Coordinator và Rebalance của Consumer Group.
  • Seek Offset: Sử dụng seekToBeginning() hoặc seek() để định vị con trỏ đọc. Đây là điểm mấu chốt để thực hiện việc replay dữ liệu.
  • Tự động xử lý Leader: Trong các phiên bản Kafka hiện đại, Kafka Client tự động quản lý metadata và định tuyến đến Leader Broker. Lập trình viên không còn phải tự viết logic tìm kiếm Leader và Replica phức tạp như trong các phiên bản 0.8.x.
  • Graceful Shutdown: Sử dụng AtomicBoolean và consumer.wakeup() để đảm bảo vòng lặp poll có thể bị phá vỡ một cách an toàn từ một luồng khác khi ứng dụng nhận tín hiệu dừng.

Thẻ: apache-kafka kafka-consumer manual-partition-assignment offset-management kafka-clients

Đăng vào ngày 9 tháng 10 lúc 22:28