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

Batching กับ Compression

producer รวม record เป็น batch และ compress ตามต้องการ ดังนั้นการรอเพิ่มอีกนิดกับ batch ที่ใหญ่ขึ้นจึงซื้อ throughput ที่สูงขึ้นมาก — เป็นแนวคิดที่ทุก client ของ Kafka ทำ แม้ปุ่มปรับจริงจะต่างกัน

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 throughput
linger.ms=10
batch.size=65536

KafkaJS ไม่มีปุ่ม linger.ms หรือ batch.size แต่จะรวม message ที่คุณส่งใน send() ครั้งเดียวเข้าด้วยกันเองแล้วส่งพร้อมกัน — คุณคุม batching ด้วยวิธีจัดกลุ่ม message เข้า send() แต่ละครั้ง (เช่น buffer record เองแล้วส่งเป็น array เดียว) ไม่ใช่ด้วย linger timer

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,
})

เพราะ 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 blocks
buffer.memory=33554432
# Total deadline covering batching + retries; must be >= request.timeout.ms + linger.ms
delivery.timeout.ms=120000
request.timeout.ms=30000

KafkaJS บรรลุเป้าหมายเดียวกัน — จำกัดว่า 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)"]
Records batch, compress, and drain within the delivery deadline

ใน 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 สะดุด

linger.ms และ batch.size ใน Java client แลกอะไรกับอะไร
compression codec ตัวไหนให้อัตราย่อดีสุดแต่กิน CPU มากสุด
KafkaJS จัดการ batching ฝั่ง producer ยังไง เทียบกับ Java client
delivery.timeout.ms (ค่า default 120000) ใน Java client คุมอะไร และ KafkaJS บรรลุเป้าหมายเดียวกันยังไง