Producer API
ไอเดียในหนึ่งประโยค
หัวข้อที่มีชื่อว่า “ไอเดียในหนึ่งประโยค”Producer สร้าง message (topic, key ที่ใส่หรือไม่ใส่ก็ได้, value) แล้วส่งเข้า send() ซึ่งคืน Promise และทำงานเสร็จแบบ asynchronous จากนั้น Kafka จะ append record นั้นลง partition ของ topic
Producer กับ message
หัวข้อที่มีชื่อว่า “Producer กับ message”สร้าง 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 clusterconst kafka = new Kafka({ clientId: 'orders-app', brokers: ['localhost:9092'] });const producer = kafka.producer();send() เป็น asynchronous — await Promise
หัวข้อที่มีชื่อว่า “send() เป็น asynchronous — await Promise”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"]