Partition ของ Apache Kafka คือแนวคิดที่อธิบายพฤติกรรมเกือบทั้งหมดของ Kafka ได้ ตั้งแต่ scale อย่างไร รับประกันลำดับข้อความแค่ไหน ไปจนถึงมี consumer ทำงานขนานกันได้กี่ตัว topic หนึ่งถูกแบ่งเป็นหลาย partition แต่ละ partition เป็น log ที่เขียนต่อท้ายอย่างเดียว (append-only) และใน consumer group หนึ่ง partition จะถูกอ่านโดย consumer ตัวเดียวเท่านั้น เข้าใจภาพนี้แล้ว เรื่องที่เหลือของ Kafka จะตามมาเอง
บทความนี้ครอบคลุมโมเดลหลัก ได้แก่ log ต่างจาก queue อย่างไร, topic, partition และ offset, key กับการเรียงลำดับ, consumer group และ retention ส่วนการดูแล broker อยู่ในบทความ รัน Kafka แบบ KRaft บน Kubernetes
Kafka คือ log ไม่ใช่ queue
ใน queue แบบดั้งเดิมอย่าง RabbitMQ พอ consumer ack ข้อความแล้ว ข้อความนั้นก็หายไป queue ทำงานเหมือนกล่องจดหมาย หยิบจดหมายออกไปแล้วก็ไม่อยู่ในกล่องอีก
Kafka เก็บข้อความเป็น log ที่เขียนต่อท้ายเรื่อย ๆ การอ่านไม่ได้ลบข้อความ ข้อความจะอยู่จนกว่า retention policy จะลบตามอายุหรือขนาด ไม่ว่าจะมีคนอ่านไปแล้วหรือยัง เปรียบได้กับสมุดบันทึก ของใหม่เขียนต่อท้าย และใครจะเปิดอ่านหน้าไหนก็ได้
ผลจากการออกแบบแบบนี้มีสามข้อ:
- Replay ได้ consumer ย้อนกลับไปประมวลผลข้อความเก่าใหม่ได้ หลังแก้ bug หรือเพิ่มฟีเจอร์
- ผู้อ่านอิสระจากกัน หลายระบบอ่าน stream เดียวกันได้ตามจังหวะของตัวเอง โดยที่ producer ไม่ต้องรู้ว่ามีใครอ่านบ้าง
- มาทีหลังก็ได้ service ใหม่เริ่มอ่านตั้งแต่ต้น log ที่ยังเก็บอยู่ แล้วสร้าง state จากประวัติได้
อีกด้านหนึ่งคือ Kafka ไม่ได้ติดตาม ack ทีละข้อความ consumer group แต่ละกลุ่มเก็บตำแหน่งแค่ค่าเดียวต่อ partition ถ้างานของคุณต้อง route ทีละข้อความ มี priority หรือต้อง delay retry ตัว queue broker มักเหมาะกว่า ดูรายละเอียดได้ใน RabbitMQ exchange และการส่งข้อความให้เชื่อถือได้
องค์ประกอบหลักของ Kafka
| คำ | ความหมาย |
|---|---|
| Topic | stream ของ event ประเภทเดียวกันที่มีชื่อ เช่น bookings หรือ payments |
| Partition | ส่วนย่อยของ topic ที่เรียงลำดับและเขียนต่อท้ายอย่างเดียว topic หนึ่งมีได้ตั้งแต่หนึ่ง partition ขึ้นไป |
| Offset | ตำแหน่งของ record ภายใน partition หนึ่ง เช่น 0, 1, 2 ไปเรื่อย ๆ |
| Producer | client ที่เขียน record เข้า topic |
| Consumer | client ที่อ่าน record จาก topic |
| Consumer group | กลุ่ม consumer ที่แบ่งงานกันอ่าน topic |
| Broker | server ของ Kafka ใน cluster มีหลายตัว และ partition กระจายอยู่บนทุกตัว |
Record หนึ่งตัวประกอบด้วย key (ใส่หรือไม่ใส่ก็ได้), value, timestamp และ header (ถ้ามี) ซึ่ง key สำคัญกว่าที่เห็นในตอนแรกมาก
Apache Kafka partition: หน่วยของการ scale และการเรียงลำดับ
Partition กระจายโหลด
Topic ที่มี 6 partition ก็คือ log อิสระ 6 ชุด แต่ละชุดกระจายอยู่บน broker ต่าง ๆ และแต่ละ partition มี leader broker หนึ่งตัวที่รับการเขียน เมื่อ leader อยู่คนละ broker การเขียนเข้า topic เดียวจึงกระจายไปทั้ง cluster ไม่ได้วิ่งผ่านเครื่องเดียว
ฝั่งอ่านก็ scale แบบเดียวกัน แต่ละ partition ให้ consumer คนละตัวอ่านขนานกันได้
การรับประกันลำดับมีแค่ภายใน partition
Kafka รับประกันลำดับ ภายใน partition เดียว ไม่ใช่ทั้ง topic record ใน partition 0 จะถูกอ่านตามลำดับที่เขียน แต่ record ใน partition 0 กับ partition 3 ไม่มีความสัมพันธ์เรื่องลำดับกันเลย
ดังนั้นถ้า event ต้องถูกประมวลผลตามลำดับ เช่น created, paid และ cancelled ของ booking เดียวกัน event พวกนั้นต้องลง partition เดียวกัน และนี่คือหน้าที่ของ key
Key เป็นตัวกำหนด partition
เมื่อ producer ส่ง record ที่มี key ตัว partitioner ค่าเริ่มต้นจะ hash key ที่ serialize แล้ว (Java client ใช้ murmur2) แล้ว mod ด้วยจำนวน partition:
partition = hash(key) % numPartitionsKey เดียวกันจะไปลง partition เดิมเสมอ ตราบใดที่จำนวน partition ไม่เปลี่ยน ใช้ business ID อย่าง bookingId หรือ userId เป็น key แล้ว event ทุกตัวของ entity นั้นจะเรียงถูกต้อง
Record ที่ไม่มี key จะกระจายไปหลาย partition ตั้งแต่ Kafka 3.3 producer จะเติม batch ให้เต็มสำหรับ partition หนึ่งก่อนค่อยขยับไป partition ถัดไป แทนการวนทีละ record ได้ batching ที่ดีกว่า แต่ไม่มีการเรียงลำดับระหว่าง record กลุ่มนี้
ผลที่ต้องวางแผนไว้ล่วงหน้ามีสามข้อ:
- เพิ่ม partition แล้ว key จะย้าย เปลี่ยนจาก 6 เป็น 12 partition ทำให้
hash(key) % nของหลาย key เปลี่ยนไป event ของ booking เดียวอาจแยกไปอยู่สอง partition ในช่วงเปลี่ยนผ่าน จึงควรเลือกจำนวน partition เผื่อไว้ตั้งแต่แรก - ลดจำนวน partition ไม่ได้ Kafka ให้เพิ่มได้อย่างเดียว ถ้าต้องการลด ต้องสร้าง topic ใหม่แล้วย้ายข้อมูล
- Hot key ทำให้เกิด hot partition ถ้าลูกค้ารายเดียวสร้าง traffic ครึ่งหนึ่งของระบบ partition ของเขาจะกลายเป็นคอขวด ไม่ว่าจะมี partition กี่ตัวก็ตาม
ส่งข้อความพร้อม key ด้วย Java
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.ACKS_CONFIG, "all"); // default since Kafka 3.0
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // default since Kafka 3.0
try (var producer = new KafkaProducer<String, String>(props)) {
var record = new ProducerRecord<>("bookings", "booking-8812", "{\"status\":\"PAID\"}");
producer.send(record, (meta, ex) -> {
if (ex != null) {
System.err.println("send failed: " + ex.getMessage());
return;
}
System.out.printf("partition=%d offset=%d%n", meta.partition(), meta.offset());
});
}ทุก record ที่มี key booking-8812 จะไปลง partition เดียวกัน event ของ booking นี้จึงถูกอ่านตามลำดับ
Offset: consumer รู้ได้อย่างไรว่าอ่านถึงไหนแล้ว
Consumer group เก็บ committed offset ไว้หนึ่งค่าต่อ partition ซึ่งก็คือตำแหน่งของ record ถัดไปที่ต้องอ่าน offset นี้เก็บอยู่ใน Kafka เอง ใน topic ภายในชื่อ __consumer_offsets ไม่ได้เก็บไว้บนเครื่องของ consumer เมื่อ consumer crash แล้วกลับมา หรือมี consumer ตัวอื่นรับ partition ไปแทน การอ่านจะต่อจาก offset ที่ commit ล่าสุด
จังหวะที่ commit เป็นตัวกำหนด delivery semantics:
- Commit หลังประมวลผล ได้ at-least-once ถ้า crash ระหว่างประมวลผลเสร็จกับ commit บาง record จะถูกประมวลผลซ้ำ handler จึงต้องเป็น idempotent
- Commit ก่อนประมวลผล ได้ at-most-once ถ้า crash อาจข้าม record ไปบางตัว
- Exactly-once ทำได้ใน pipeline ที่อ่านจาก Kafka แล้วเขียนกลับ Kafka ด้วย transaction แต่พอมี database ภายนอกเข้ามาเกี่ยว ก็ต้องกลับไปพึ่งการเขียนแบบ idempotent
auto.offset.reset กำหนดว่า group จะเริ่มอ่านจากไหนเมื่อยังไม่มี committed offset ค่า earliest อ่านตั้งแต่ต้น log ที่ยังเก็บอยู่ ส่วน latest ซึ่งเป็นค่าเริ่มต้น อ่านเฉพาะ record ใหม่
ตัวอย่าง consumer ที่ commit หลังประมวลผล:
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "analytics");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
try (var consumer = new KafkaConsumer<String, String>(props)) {
consumer.subscribe(List.of("bookings"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
System.out.printf("p=%d offset=%d key=%s%n", r.partition(), r.offset(), r.key());
}
consumer.commitSync(); // commit only after the batch is processed
}
}การ replay ก็คือการ reset offset นั่นเอง หยุด consumer ของ group นั้นก่อน แล้วรัน:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group analytics --topic bookings \
--reset-offsets --to-datetime 2026-09-01T00:00:00.000 --executeKafka consumer group: กฎสองข้อ
กฎข้อ 1: consumer ที่ใช้ group ID เดียวกันจะแบ่งงานกัน Kafka มอบแต่ละ partition ให้ consumer ในกลุ่มแค่ตัวเดียว ถ้ามี consumer 3 ตัวใน group analytics อ่าน topic ที่มี 6 partition แต่ละตัวจะได้ 2 partition และแต่ละ record ถูกประมวลผลแค่ครั้งเดียวในกลุ่ม
กฎข้อ 2: คนละ group ID ต่างคนต่างได้ครบทุกข้อความ group analytics กับ group email อ่านทุก record ใน bookings แยกกันอิสระ แต่ละกลุ่มมี offset ของตัวเอง event การจองหนึ่งตัวจึงถูกประมวลผลหนึ่งครั้งโดย analytics และอีกหนึ่งครั้งโดย email
สรุปสั้น ๆ คือ topic บอกว่ามีข้อมูล อะไร ส่วน group บอกว่า ใคร อ่าน และแบ่งงานกันอย่างไร
จำนวน partition คือเพดานของ consumer ที่ทำงานขนานได้
เพราะหนึ่ง partition ไปอยู่กับ consumer ได้ตัวเดียวต่อ group จำนวน partition จึงเป็นเพดานของการทำงานขนานในกลุ่มนั้น ถ้ามี 3 partition consumer ตัวที่ 4 ใน group เดียวกันจะว่างงาน มันยังช่วยรับงานแทนได้ถ้าตัวอื่นตาย แต่ไม่ได้เพิ่ม throughput
เมื่อมี consumer เข้าหรือออก หรือมีตัวไหนไม่ poll นานเกิน max.poll.interval.ms กลุ่มจะ rebalance และจัดสรร partition ใหม่ การ rebalance บ่อย ๆ เป็นสาเหตุของ lag ที่พบบ่อย ควรจำกัดปริมาณงานต่อการ poll แต่ละครั้ง และตั้ง CooperativeStickyAssignor เพื่อให้หยุดชั่วคราวเฉพาะ partition ที่ย้ายจริง นอกจากนี้ Kafka 4.0 ยังเปิดใช้ consumer group protocol ใหม่ (group.protocol=consumer) เป็น GA ซึ่งย้ายการจัดสรรไปทำที่ broker และลดผลกระทบจาก rebalance ลงอีก
ดูสถานะและ lag ของ group:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group analyticsคอลัมน์ LAG บอกว่า committed offset ของแต่ละ partition ตามหลังท้าย log อยู่เท่าไร ถ้า lag โตอยู่ partition เดียว มักเป็นสัญญาณของ hot key
Retention และ log compaction
Kafka ลบข้อมูลตาม policy ไม่ใช่เพราะมีคนอ่านไปแล้ว:
retention.msกำหนดระยะเวลาเก็บ record ค่าเริ่มต้นของ broker คือ 7 วันretention.bytesจำกัดขนาด ต่อ partition ค่าเริ่มต้นคือไม่จำกัด- การลบทำทีละ segment file ทั้งไฟล์ record จึงอาจอยู่นานกว่าค่าที่ตั้งไว้เล็กน้อย
ถ้าตั้ง cleanup.policy=compact Kafka จะเก็บ record ล่าสุดของแต่ละ key ไว้อย่างน้อยหนึ่งตัว และลบตัวเก่าออก topic จะกลายเป็น changelog ของ state ปัจจุบัน เช่น profile ล่าสุดของ user แต่ละคน ส่วน record ที่ value เป็น null (เรียกว่า tombstone) ใช้บอกให้ลบ key นั้น
สร้าง topic พร้อมตั้งค่าให้ชัดเจน (ตัวอย่างนี้เก็บ 14 วัน):
kafka-topics.sh --bootstrap-server localhost:9092 --create \
--topic bookings --partitions 12 --replication-factor 3 \
--config retention.ms=1209600000 --config min.insync.replicas=2ใช้ Kafka เมื่อไหร่ถึงจะคุ้ม
Kafka เหมาะเมื่อ:
- ต้องการ event stream ทุกการจอง การค้นหา และการคลิก ถูกบันทึกเป็น event ให้ service อื่นตอบสนอง
- หลายระบบใช้ event ชุดเดียวกัน analytics, ระบบอีเมล, fraud detection และ data warehouse ต่างอ่านด้วย consumer group ของตัวเอง การส่งแบบ realtime ก็เป็น consumer อีกตัวได้ เหมือนใน การ scale WebSocket ด้วย pub/sub
- ต้อง replay ได้ ประมวลผลประวัติใหม่หลังแก้ bug หรือสร้าง service ใหม่จาก event ในอดีต นี่คือเหตุผลที่ Kafka มักอยู่ใต้สถาปัตยกรรม event sourcing และ CQRS
- Throughput สูง การเขียนต่อท้ายแบบ sequential, การรวม batch ฝั่ง producer และ zero-copy จาก disk ไป network ทำให้ Kafka รับปริมาณที่ broker แบบทีละข้อความรับไม่ไหว
Kafka ไม่ค่อยเหมาะกับงาน request/response, งานที่ต้องมี priority หรือ delay ทีละข้อความ และระบบเล็กที่ queue ตัวเดียวก็พอ การดูแล Kafka cluster มีต้นทุนจริง
คำถามที่พบบ่อย
Topic ใน Kafka ควรมีกี่ partition
เริ่มจากระดับการทำงานขนานที่ต้องการ อย่างน้อยควรเท่ากับจำนวน consumer สูงสุดที่คาดว่าจะมีใน group เดียว บวกเผื่อไว้ เพราะการเพิ่ม partition ทีหลังทำให้ key ย้าย แต่จำนวนที่สูงมากก็เพิ่มภาระให้ broker และ client อย่าตั้งเป็นหลักร้อยโดยไม่มีเหตุผล
Kafka รับประกันลำดับข้อความไหม
รับประกันแค่ภายใน partition เดียว ให้ record ที่เกี่ยวข้องกันใช้ key เดียวกันเพื่อให้ลง partition เดียวกัน
ถ้ามี consumer มากกว่า partition จะเกิดอะไรขึ้น
Consumer ส่วนที่เกินใน group นั้นจะไม่ได้ partition และว่างงานจนกว่าจะมี rebalance ที่มอบ partition ให้
Kafka ลบข้อความหลังถูกอ่านไหม
ไม่ลบ record ถูกลบตามเวลาหรือขนาดของ retention หรือจาก compaction ไม่เกี่ยวกับว่าถูกอ่านไปแล้วหรือยัง
Topic กับ consumer group ต่างกันอย่างไร
Topic คือที่ที่ข้อมูลถูกเขียนเข้าไป ส่วน consumer group คือกลุ่มผู้อ่านที่แบ่ง partition ของ topic กันอ่าน และเก็บ offset ของตัวเอง
สรุปสิ่งที่ควรจำ
- Kafka คือ log การอ่านไม่ได้ลบข้อความ replay และการมีผู้อ่านอิสระหลายรายจึงมาพร้อมในตัว
- Partition เป็นหน่วยของทั้งการ scale และการเรียงลำดับ เลือกจำนวนโดยเผื่อไว้
- อะไรที่ต้องเรียงลำดับ ให้ใช้ business ID เป็น key
- Group ID เดียวกันแบ่งงานกัน คนละ group ID ได้ครบทุกข้อความ
- Commit หลังประมวลผล และทำ handler ให้เป็น idempotent
ขั้นต่อไป ดูว่า broker, controller และ replication ทำงานร่วมกันอย่างไรใน รัน Kafka แบบ KRaft บน Kubernetes ถ้ากำลังตัดสินใจว่า Kafka เหมาะกับระบบที่กำลังสร้างหรือไม่ ทีม Vectorkub ช่วยออกแบบได้
