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

Consumer Groups

Consumer group ช่วยให้หลาย worker (consumer) อ่านจาก stream เดียวกันได้ โดยที่แต่ละ entry จะถูกส่งไปยัง consumer เพียงหนึ่งคนเท่านั้นในกลุ่ม Redis ติดตามว่า entry ใดถูกส่งแล้วแต่ยังไม่ได้รับการ acknowledge ซึ่งทำให้เกิด at-least-once processing ได้

127.0.0.1:6379> XADD jobs * type "email" to "[email protected]"
"1700000000000-0"
127.0.0.1:6379> XADD jobs * type "sms" to "+6681234567"
"1700000000001-0"
127.0.0.1:6379> XGROUP CREATE jobs workers $ MKSTREAM
OK

$ หมายถึง “เริ่มจาก entry ล่าสุด” — เฉพาะ message ที่เพิ่มใหม่หลังจากสร้าง group เท่านั้นที่จะถูกส่ง ใช้ 0 เพื่ออ่านตั้งแต่ต้น stream MKSTREAM จะสร้าง stream ให้อัตโนมัติหาก stream นั้นยังไม่มีอยู่

127.0.0.1:6379> XADD jobs * type "email" to "[email protected]"
"1700000000002-0"
127.0.0.1:6379> XREADGROUP GROUP workers consumer-1 COUNT 1 STREAMS jobs >
1) 1) "jobs"
2) 1) 1) "1700000000002-0"
2) 1) "type"
2) "email"
3) "to"

ID พิเศษ > หมายความว่า “ส่ง entry ใหม่ที่ยังไม่ได้รับมาให้ฉัน” เมื่อส่งแล้ว entry จะถูกวางไว้ใน Pending Entries List (PEL) ของ consumer-1 จนกว่าจะได้รับการ acknowledge อย่างชัดเจน

127.0.0.1:6379> XACK jobs workers 1700000000002-0
(integer) 1

เมื่อ acknowledge แล้ว entry จะถูกลบออกจาก PEL เรียก XACK เฉพาะหลังจากที่ worker ประมวลผล message สำเร็จแล้วเท่านั้น — การ acknowledge เร็วเกินไปอาจทำให้ข้อมูลสูญหายหาก worker crash ระหว่างประมวลผล

127.0.0.1:6379> XPENDING jobs workers - + 10
1) 1) "1700000000002-0"
2) "consumer-1"
3) (integer) 4321
4) (integer) 1

แต่ละแถวแสดง: ID ของ entry, consumer ที่ถือครอง entry นั้น, ระยะเวลาที่ idle เป็นมิลลิวินาที และจำนวนครั้งที่ส่งแล้ว เครื่องมือนี้ใช้ตรวจหา consumer ที่ค้างหรือทำงานช้า

XADD jobs * type "email" to "[email protected]"
XADD jobs * type "sms" to "+6681234567"
XGROUP CREATE jobs workers 0 MKSTREAM
XREADGROUP GROUP workers consumer-1 COUNT 10 STREAMS jobs >
XACK jobs workers 1700000000000-0
XPENDING jobs workers - + 10
ตัวเลือกBenefitCost
Consumer group แบ่งงานให้หลาย workerscale การประมวลผลแนวนอนได้ แต่ละ entry ไปหา consumer เดียว ไม่ซ้ำงานต้องจัดการ XACK เองทุก entry เพิ่มความซับซ้อนเทียบกับ XREAD ธรรมดา
At-least-once delivery ผ่าน PELไม่มี entry สูญหายแม้ consumer crash กลางคันconsumer ต้องเขียน logic ให้ idempotent เพราะ entry อาจถูกส่งซ้ำได้เมื่อ claim ใหม่
  • ไม่เรียก XACK หลังประมวลผลสำเร็จ — entry จะค้างอยู่ใน PEL ตลอดไป ทำให้ PEL โตขึ้นเรื่อย ๆ และมองไม่เห็นว่า entry ไหนจบงานแล้วจริง ๆ
  • XACK ก่อนประมวลผลเสร็จ — ถ้า worker crash หลัง ack แต่ก่อนทำงานจริงเสร็จ entry นั้นจะหายไปจากระบบ tracking ทั้งที่งานยังไม่เสร็จ ควร ack หลังบันทึกผลสำเร็จเท่านั้น
  • ไม่มี process คอย claim entry ที่ค้างจาก consumer ที่ตายไปแล้ว — ถ้าไม่ใช้ XCLAIM/XAUTOCLAIM ตรวจสอบ PEL เป็นระยะ entry ของ consumer ที่ crash จะไม่มีวันถูกประมวลผลต่อ

💡 ตัวอย่างจากของจริง

Sidekiq-style background job queue — ระบบประมวลผลงานเบื้องหลัง (ส่งอีเมล, resize รูป) ใช้ consumer group กระจายงานให้ worker หลายตัวพร้อม guarantee ว่างานจะไม่หายแม้ worker ตัวหนึ่ง crash

Log/event processing แบบกระจาย — ทีมข้อมูลใช้ consumer group หลายตัวอ่าน stream เดียวกันเพื่อประมวลผลแบบขนาน โดยอาศัย PEL และ XCLAIM รับประกันว่าไม่มี event ไหนถูกข้ามไปแม้ instance ล้มเหลวระหว่างทาง

`XREADGROUP GROUP g consumer STREAMS s >` หมายความว่าอะไร?
เกิดอะไรขึ้นกับ stream entry หลังจากเรียก XACK?
XPENDING แสดงอะไร?
ใน `XGROUP CREATE jobs workers 0 MKSTREAM` ตัวเลข `0` หมายความว่าอะไร?