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

Consumer Groups & Rebalancing

consumer group คือกลุ่มของ consumer ที่แชร์ partition ของ topic กัน — แต่ละ partition ถูกจัดการโดย consumer เพียงตัวเดียวใน group — และ rebalancing คือกระบวนการจัดสรร partition ใหม่เมื่อสมาชิกใน group เปลี่ยนไป

consumer ทุกตัวประกาศ groupId Kafka จะยก partition แต่ละอันให้ consumer เพียง ตัวเดียว ภายใน group ดังนั้น topic ที่มี 6 partition จะถูกประมวลผลพร้อมกันได้สูงสุด 6 consumer เพิ่ม consumer ตัวที่ 7 เข้าไปก็จะนั่งว่าง — เพดานของ parallelism คือจำนวน partition ไม่ใช่จำนวน consumer

group ต่างกันเป็นอิสระต่อกัน: analytics กับ billing ต่างได้ stream เต็ม ๆ ของตัวเอง และ track offset ของตัวเอง ส่วนภายใน group เดียวกัน partition จะถูกแบ่งกัน

คุณตรวจดู group, สมาชิก และ lag ได้จาก CLI:

Terminal window
# Show members, assigned partitions, and lag for a group
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group billing
flowchart TB
  subgraph t["Topic orders: 3 partitions"]
    p0["orders-0"]
    p1["orders-1"]
    p2["orders-2"]
  end
  subgraph g1["Group billing"]
    c1["consumer A"]
    c2["consumer B"]
  end
  p0 --> c1
  p1 --> c1
  p2 --> c2
  t --> g2["Group analytics (own offsets, full copy)"]
One group divides partitions; another group reads the same topic independently

เมื่อ consumer เข้า, ออก หรือตาย — หรือมีการเพิ่ม partition — group ต้อง rebalance: กระจาย partition ใหม่ไปยังสมาชิกที่มีอยู่ตอนนั้น protocol แบบ classic ทำสิ่งนี้ด้วยขั้นตอน stop-the-world: consumer ทุกตัวหยุดประมวลผล คืน partition ทั้งหมด แล้วรอ assignment ใหม่ บน group ขนาดใหญ่ การหยุดแบบนั้นเจ็บมาก

Kafka 4.0 มาพร้อม next-generation consumer rebalance protocol (KIP-848) ในสถานะ generally available สำหรับ Java client โดยย้าย logic การ assign ไปไว้ที่ group coordinator ฝั่ง broker และทำให้ rebalance เป็นแบบ incremental — เพิกถอนเฉพาะ partition ที่ต้องย้ายจริง ๆ ส่วนที่เหลือประมวลผลต่อได้ ผลลัพธ์คือ rebalance สั้นลงและรบกวนงานน้อยลง โดยเฉพาะกับ group ขนาดใหญ่

# Java-client opt-in: the new consumer group protocol (KIP-848), GA in Kafka 4.0
group.protocol=consumer
# The classic (older) protocol is still available:
# group.protocol=classic

consumer ของ KafkaJS — ตัวที่ตัวอย่างในคอร์สนี้ใช้ — ตอนนี้พูด protocol แบบ classic เท่านั้น ไม่มี setting group.protocol ให้เปลี่ยน แต่นั่นไม่ได้แปลว่าเสียประโยชน์ไปทั้งหมด เพราะการปรับปรุงของ KIP-848 ส่วนใหญ่อยู่ที่ group coordinator ฝั่ง broker ดังนั้น cluster KRaft 4.x ก็ยังได้ประโยชน์จาก coordinator ตัวใหม่ที่มีประสิทธิภาพกว่า แม้จะคุยกับ consumer แบบ classic-protocol อย่าง KafkaJS การรองรับ protocol ใหม่บน JavaScript client (เช่น @confluentinc/kafka-javascript ของ Confluent) จะตามมาเมื่อมีการ implement ที่ต้นทาง

ภายใน consumer group เดียว มี consumer กี่ตัวที่อ่าน partition หนึ่ง ๆ?
topic มี 4 partition แล้วคุณสตาร์ท 6 consumer ใน group เดียว จะเกิดอะไรขึ้น?
อะไรดีขึ้นเรื่อง rebalancing ใน Kafka 4.0 ด้วย KIP-848 และครอบคลุม consumer ของ KafkaJS วันนี้ไหม?
consumer group สองกลุ่มที่อ่าน topic เดียวกันสัมพันธ์กันอย่างไร?