Bài 5: Producer API
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:
bootstrap.servers: Địa chỉ của một hoặc nhiều Kafka broker (ví dụ:
localhost:9092).key.serializer: Lớp dùng để serialize khóa của message (ví dụ:
org.apache.kafka.common.serialization.StringSerializer).value.serializer: Lớp dùng để serialize giá trị của message (ví dụ:
org.apache.kafka.common.serialization.StringSerializer).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
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>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
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.
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-beginningBạ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.