Skip to main content

Command Palette

Search for a command to run...

Bài 5: Producer API

Published
3 min readView as Markdown

Mục tiêu

Hiểu và sử dụng Producer API để gửi thông điệp vào Kafka.

Nội dung chính

Khái niệm về Producer trong Kafka

Producer là các ứng dụng hoặc dịch vụ gửi message vào các topic trong Kafka. Producer chịu trách nhiệm lựa chọn partition để gửi message và đảm bảo message được gửi đi đúng cách. Các producer có thể gửi message đồng bộ hoặc không đồng bộ.

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

Khi cấu hình một Kafka producer, 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. key.serializer: Lớp dùng để serialize khóa của message (ví dụ: org.apache.kafka.common.serialization.StringSerializer).

  3. value.serializer: Lớp dùng để serialize giá trị của message (ví dụ: org.apache.kafka.common.serialization.StringSerializer).

  4. acks: Xác định mức độ xác nhận cần thiết từ Kafka trước khi coi message là đã gửi thành công (ví dụ: 0, 1, all).

Gửi thông điệp đồng bộ và không đồng bộ

  • Gửi đồng bộ: Producer chờ đợi một phản hồi từ Kafka broker trước khi tiếp tục gửi message tiếp theo.

  • Gửi không đồng bộ: Producer gửi message mà không chờ đợi phản hồi, sử dụng callback để xử lý kết quả khi có phản hồi từ broker.

Thực hành

Viết một ứng dụng Java sử dụng Producer API để gửi thông điệp vào 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-producer-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 Producer

     import org.apache.kafka.clients.producer.KafkaProducer;
     import org.apache.kafka.clients.producer.ProducerConfig;
     import org.apache.kafka.clients.producer.ProducerRecord;
     import org.apache.kafka.clients.producer.RecordMetadata;
     import org.apache.kafka.common.serialization.StringSerializer;
    
     import java.util.Properties;
     import java.util.concurrent.ExecutionException;
    
     public class KafkaProducerExample {
         public static void main(String[] args) {
             Properties properties = new Properties();
             properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
             properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
             properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
             properties.put(ProducerConfig.ACKS_CONFIG, "all");
    
             KafkaProducer<String, String> producer = new KafkaProducer<>(properties);
    
             String topic = "my-topic";
             String key = "my-key";
             String value = "Hello, Kafka!";
    
             // Gửi đồng bộ
             try {
                 RecordMetadata metadata = producer.send(new ProducerRecord<>(topic, key, value)).get();
                 System.out.printf("Sent message to topic %s partition %d offset %d%n", metadata.topic(), metadata.partition(), metadata.offset());
             } catch (InterruptedException | ExecutionException e) {
                 e.printStackTrace();
             }
    
             // Gửi không đồng bộ
             producer.send(new ProducerRecord<>(topic, key, value), (metadata, exception) -> {
                 if (exception == null) {
                     System.out.printf("Sent message to topic %s partition %d offset %d%n", metadata.topic(), metadata.partition(), metadata.offset());
                 } else {
                     exception.printStackTrace();
                 }
             });
    
             producer.close();
         }
     }
    

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

  1. Chạy Producer:

    • Chạy chương trình Java trên và kiểm tra đầu ra để đảm bảo message đã được gửi thành công.
  2. Sử dụng Kafka Console Consumer để xác minh:

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

        bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic my-topic --from-beginning
      
    • Bạn sẽ thấy message "Hello, Kafka!" được nhận từ topic.

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

More from this blog

hoangkim

366 posts