Consumer Groups
Consumer Group คืออะไร
หัวข้อที่มีชื่อว่า “Consumer Group คืออะไร”Consumer group ช่วยให้หลาย worker (consumer) อ่านจาก stream เดียวกันได้ โดยที่แต่ละ entry จะถูกส่งไปยัง consumer เพียงหนึ่งคนเท่านั้นในกลุ่ม Redis ติดตามว่า entry ใดถูกส่งแล้วแต่ยังไม่ได้รับการ acknowledge ซึ่งทำให้เกิด at-least-once processing ได้
สร้าง Group — XGROUP CREATE
หัวข้อที่มีชื่อว่า “สร้าง Group — XGROUP CREATE”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 $ MKSTREAMOK$ หมายถึง “เริ่มจาก entry ล่าสุด” — เฉพาะ message ที่เพิ่มใหม่หลังจากสร้าง group เท่านั้นที่จะถูกส่ง ใช้ 0 เพื่ออ่านตั้งแต่ต้น stream MKSTREAM จะสร้าง stream ให้อัตโนมัติหาก stream นั้นยังไม่มีอยู่
XREADGROUP — อ่านในฐานะ consumer
หัวข้อที่มีชื่อว่า “XREADGROUP — อ่านในฐานะ consumer”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" 4) "[email protected]"ID พิเศษ > หมายความว่า “ส่ง entry ใหม่ที่ยังไม่ได้รับมาให้ฉัน” เมื่อส่งแล้ว entry จะถูกวางไว้ใน Pending Entries List (PEL) ของ consumer-1 จนกว่าจะได้รับการ acknowledge อย่างชัดเจน
XACK — ยืนยันการประมวลผล
หัวข้อที่มีชื่อว่า “XACK — ยืนยันการประมวลผล”127.0.0.1:6379> XACK jobs workers 1700000000002-0(integer) 1เมื่อ acknowledge แล้ว entry จะถูกลบออกจาก PEL เรียก XACK เฉพาะหลังจากที่ worker ประมวลผล message สำเร็จแล้วเท่านั้น — การ acknowledge เร็วเกินไปอาจทำให้ข้อมูลสูญหายหาก worker crash ระหว่างประมวลผล
XPENDING — ตรวจสอบ Entry ที่ยังไม่ได้ Acknowledge
หัวข้อที่มีชื่อว่า “XPENDING — ตรวจสอบ Entry ที่ยังไม่ได้ Acknowledge”127.0.0.1:6379> XPENDING jobs workers - + 101) 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ข้อแลกเปลี่ยน
หัวข้อที่มีชื่อว่า “ข้อแลกเปลี่ยน”| ตัวเลือก | Benefit | Cost |
|---|---|---|
| Consumer group แบ่งงานให้หลาย worker | scale การประมวลผลแนวนอนได้ แต่ละ 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 ล้มเหลวระหว่างทาง