Batching กับ Compression
ไอเดียในหนึ่งประโยค
หัวข้อที่มีชื่อว่า “ไอเดียในหนึ่งประโยค”producer รวม record เป็น batch และ compress ตามต้องการ ดังนั้นการรอเพิ่มอีกนิดกับ batch ที่ใหญ่ขึ้นจึงซื้อ throughput ที่สูงขึ้นมาก — เป็นแนวคิดที่ทุก client ของ Kafka ทำ แม้ปุ่มปรับจริงจะต่างกัน
linger.ms กับ batch.size แลก latency กับ throughput
หัวข้อที่มีชื่อว่า “linger.ms กับ batch.size แลก latency กับ throughput”producer ไม่ได้ส่งทีละ record แต่สะสม record ต่อ partition เข้าเป็น batch แล้วส่ง batch เป็น request เดียว มีสองค่ามาตรฐานที่กำหนดรูปร่างของ batch — นี่คือ ชื่อ config ของ Java client / librdkafka ที่คุณจะเจอใน official docs และ client ส่วนใหญ่:
batch.size— ขนาดสูงสุดเป็น byte ของ batch เดียว batch จะถูกส่งเมื่อเต็มถึงขนาดนี้ (หรือเมื่อlinger.msหมดเวลา)linger.ms— producer จะ รอ record เพิ่มนานแค่ไหนก่อนส่ง batch ที่ยังไม่เต็ม ค่า default คือ0(ส่งให้เร็วที่สุด) เพิ่มเป็น เช่น10จะให้ record สะสมได้มากขึ้น เกิด batch ที่เต็มกว่า
batch ที่ใหญ่ขึ้นหมายถึง request น้อยลง compress ได้ดีขึ้น และ throughput สูงขึ้น — แลกกับ latency ที่เพิ่มไม่กี่ millisecond ต่อ record นี่คือคันโยกหลัก
# Wait up to 10 ms to fill a batch, and allow bigger batches for throughputlinger.ms=10batch.size=65536KafkaJS ไม่มีปุ่ม linger.ms หรือ batch.size แต่จะรวม message ที่คุณส่งใน send() ครั้งเดียวเข้าด้วยกันเองแล้วส่งพร้อมกัน — คุณคุม batching ด้วยวิธีจัดกลุ่ม message เข้า send() แต่ละครั้ง (เช่น buffer record เองแล้วส่งเป็น array เดียว) ไม่ใช่ด้วย linger timer
Compress batch
หัวข้อที่มีชื่อว่า “Compress batch”compression ถูกทำ ต่อ batch ดังนั้น batch ที่ใหญ่กว่าจึง compress ได้ดีกว่า ทุก client เลือก codec บนเส้นความเร็วต่ออัตราการย่อเดียวกัน:
lz4— เร็วมาก อัตราย่อดี เป็นตัวเลือก default ที่นิยมzstd— อัตราย่อเยี่ยม ยังเร็วอยู่ เหมาะเมื่อ network หรือ storage เป็นคอขวดsnappy— เร็ว อัตราย่อพอประมาณgzip— อัตราย่อดีสุด แต่กิน CPU มากสุดและช้าสุดnone— ไม่ compress (ค่า default)
ใน Java client นี่คือ property compression.type ใน KafkaJS compression เป็น option ต่อการ send ที่ใช้ CompressionTypes:
import { Kafka, CompressionTypes } from 'kafkajs'
// KafkaJS ships GZIP natively; snappy/lz4/zstd need an extra codec// package registered via CompressionCodecs before you can use them.await producer.send({ topic: 'orders', compression: CompressionTypes.GZIP, messages,})buffer.memory, back-pressure และเส้นตายของการส่ง
หัวข้อที่มีชื่อว่า “buffer.memory, back-pressure และเส้นตายของการส่ง”เพราะ send() เป็น asynchronous record จึงนั่งอยู่ใน buffer ในหน่วยความจำก่อน I/O layer จะส่งออก ใน Java client buffer.memory จำกัดขนาดรวมของ buffer นั้น เมื่อ buffer เต็ม send() จะ block (นานสุดตาม max.block.ms) เป็นการใส่ back-pressure เพื่อไม่ให้ producer ที่เร็วเกินไปกินหน่วยความจำจนหมด KafkaJS ไม่มี setting buffer แบบนี้ให้ปรับ ไม่ได้ pre-buffer record ก่อน send() แบบที่ Java client ทำ จึงไม่มีคู่ buffer.memory/max.block.ms ให้จูน
Java client ยังคุมทั้งเส้นทางด้วยเส้นตาย: delivery.timeout.ms (ค่า default 120000) คือ เส้นตายรวม ตั้งแต่ตอน send() คืนค่าจนถึงตอน record ถูก acknowledge — ครอบคลุมทั้ง delay ของ batching, การส่งผ่าน network และ retry ทั้งหมด ควรตั้งให้ >= request.timeout.ms + linger.ms ถ้าส่ง record ไม่สำเร็จภายในเวลานี้ การ send จะ fail ถาวร
# 32 MiB send buffer; when full, send() applies back-pressure and blocksbuffer.memory=33554432# Total deadline covering batching + retries; must be >= request.timeout.ms + linger.msdelivery.timeout.ms=120000request.timeout.ms=30000KafkaJS บรรลุเป้าหมายเดียวกัน — จำกัดว่า send จะ retry ได้นานแค่ไหน — ด้วยวิธีต่างออกไป: ผ่าน retry policy ระดับ client (retries, initialRetryTime, maxRetryTime) ซึ่งอยู่ใน Error Handling & Retries ไม่มี deadline เดี่ยว ๆ retry จะหยุดเองเมื่อจำนวนครั้งใน policy หมด
flowchart LR send["send() adds record to buffer"] --> buf["'buffer.memory': full means back-pressure"] buf --> batch["Batch fills to 'batch.size' or 'linger.ms' elapses"] batch --> comp["Compress batch: lz4 / zstd / snappy / gzip"] comp --> io["I/O thread sends; retries allowed"] io --> deadline["Must finish within 'delivery.timeout.ms' (120000)"]
จูนเพื่อ throughput
หัวข้อที่มีชื่อว่า “จูนเพื่อ throughput”ใน Java client ถ้าจะดัน throughput ขึ้น เพิ่ม linger.ms อีกนิด (5-20 ms) ขยาย batch.size และเปิด compression แบบ lz4 หรือ zstd — สามอย่างนี้ทำงานร่วมกัน เพราะรอมากขึ้นทำให้ batch เต็มกว่า และ batch ที่เต็มกว่าก็ compress ได้ดีกว่า เผื่อ buffer.memory ให้พอสบายเพื่อไม่ให้ back-pressure มาสะดุด และทำให้ delivery.timeout.ms ยาวพอจะรองรับ retry ตอน broker สะดุด ใน KafkaJS เพราะไม่มีปุ่ม batching ให้จูน ให้เน้นงาน throughput ที่: ส่ง message array ที่ใหญ่ขึ้นต่อการ send() หนึ่งครั้ง เลือก compression: CompressionTypes.GZIP (หรือ codec LZ4/ZSTD ที่ลงทะเบียนไว้) และตั้ง retry policy ให้กว้างพอจะทนต่อ broker สะดุด