ข้ามไปยังเนื้อหา

Producer API

Producer สร้าง message (topic, key ที่ใส่หรือไม่ใส่ก็ได้, value) แล้วส่งเข้า send() ซึ่งคืน Promise และทำงานเสร็จแบบ asynchronous จากนั้น Kafka จะ append record นั้นลง partition ของ topic

สร้าง producer ครั้งเดียวจาก Kafka client, connect() ครั้งเดียว แล้วใช้ซ้ำตลอดอายุแอปได้เลย เราไม่เขียน byte ดิบไปหา broker เอง แต่ส่ง topic พร้อม message หนึ่งหรือหลายอันให้ send()

message พก 3 อย่างที่เราสนใจมากสุด:

  • topic — ชื่อ stream ที่จะเขียนลง (ระบุตอนเรียก send(), จำเป็น)
  • key — ใส่หรือไม่ก็ได้ record ที่มี key เดียวกันจะไปลง partition เดียวกัน นั่นคือวิธีที่ Kafka รักษา ordering ต่อ key
  • value — payload หรือตัว event data นั่นเอง

เพราะ broker เก็บแค่ byte เท่านั้น key และ value ของ message จึงต้องเป็น string, Buffer หรือ null KafkaJS ไม่มี serializer class แบบ pluggable — คุณ serialize payload เอง ปกติใช้ JSON.stringify(...) สำหรับ value ที่เป็น structured และใช้ string ธรรมดาสำหรับ key (client อื่นอย่าง Java จะตั้งค่า key.serializer และ value.serializer เป็น class แทน)

import { Kafka } from 'kafkajs';
// brokers: the initial broker list the client contacts to discover the whole cluster
const kafka = new Kafka({ clientId: 'orders-app', brokers: ['localhost:9092'] });
const producer = kafka.producer();

send() ไม่ block รอจน record ถึง broker แต่ KafkaJS จะ buffer และ batch ให้ภายใน แล้วคืน Promise กลับมา ถ้าอยากรู้ผลลัพธ์ — partition และ offset ที่ Kafka assign ให้ — ก็ await promise นั้น (หรือใช้ .then()) และถ้าล้มเหลว promise จะ throw ดังนั้นครอบด้วย try/catch

import { Kafka } from 'kafkajs';
const kafka = new Kafka({ clientId: 'orders-app', brokers: ['localhost:9092'] });
const producer = kafka.producer();
async function main() {
await producer.connect(); // connect once, then reuse the producer
try {
// one message: topic "orders", key "order-42", JSON-serialized value
const [meta] = await producer.send({
topic: 'orders',
messages: [{ key: 'order-42', value: JSON.stringify({ amount: 19.99 }) }],
});
// success — Kafka assigned this record a partition and a base offset
console.log(`sent to ${meta.topicName}-${meta.partition} @ offset ${meta.baseOffset}`);
} catch (err) {
// the promise rejected — log, retry, or route to a dead-letter path
console.error('send failed:', err);
} finally {
await producer.disconnect(); // flush buffered records, release resources
}
}
main();
flowchart LR
  app["Your app builds a message"] --> ser["Serialize key/value yourself, e.g. JSON.stringify"]
  ser --> buf["Producer buffer (per partition)"]
  buf --> io["KafkaJS batches and sends"]
  io --> broker["Broker appends to partition"]
  broker --> cb["Promise resolves: metadata, or throws"]
จาก message สู่ offset ที่ถูก assign
message ของ Kafka มี 3 field หลักอะไรบ้าง
send() คืนค่าอะไร และทำงานเสร็จเมื่อไร
key และ value ของ message ถูกแทนด้วยอะไรใน KafkaJS
brokers array บน new Kafka({...}) ใช้ทำอะไร