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

Redis Streams

Redis Stream คือโครงสร้างข้อมูลที่ทำงานเหมือน log ที่มีลำดับและถาวร ต่างจาก Pub/Sub ที่ข้อความจะหายไปทันทีที่ส่งถึงผู้รับ — stream จะ เก็บทุก entry ไว้และให้คุณอ่านซ้ำได้ทุกเมื่อ จากตำแหน่งใดก็ได้

แต่ละ entry มีสองส่วน:

  • ID — identifier เฉพาะที่อิงตาม timestamp ในรูปแบบ <milliseconds>-<sequence> คุณสามารถให้ Redis สร้างให้อัตโนมัติโดยส่ง * หรือจะกำหนดเองก็ได้
  • Fields — คู่ key-value หนึ่งคู่หรือมากกว่า ที่เป็น payload ของ entry คล้ายกับ hash

Streams เหมาะสำหรับ event log, activity feed, metrics pipeline และงานใด ๆ ที่ต้องการทั้งการส่งแบบ real-time และการ replay ประวัติย้อนหลัง

XADD key id field value [field value ...] เพิ่ม entry ใหม่เข้าไปใน stream การใช้ * เป็น ID จะให้ Redis กำหนด ID แบบ timestamp ที่เพิ่มขึ้นเรื่อย ๆ ให้อัตโนมัติ

127.0.0.1:6379> XADD events * action "login" user "ada"
"1526919030474-0"
127.0.0.1:6379> XADD events * action "purchase" user "ada" item "book"
"1526919030474-1"
127.0.0.1:6379> XLEN events
(integer) 2

XLEN คืนจำนวน entry ที่มีอยู่ใน stream ขณะนั้น — เหมือนกับ LLEN สำหรับ list ทุกประการ

XRANGE key start end [COUNT n] คืน entry จากเก่าสุดไปใหม่สุด โดยมี ID พิเศษสองตัวที่ช่วยให้ query ง่ายขึ้น:

  • - — entry แรกสุดใน stream
  • + — entry ล่าสุดใน stream

XREVRANGE key end start กลับทิศทาง คืน entry จากใหม่สุดไปเก่าสุด

127.0.0.1:6379> XRANGE events - +
1) 1) "1526919030474-0"
2) 1) "action"
2) "login"
3) "user"
4) "ada"
2) 1) "1526919030474-1"
2) 1) "action"
2) "purchase"
3) "user"
4) "ada"
5) "item"
6) "book"
127.0.0.1:6379> XRANGE events - + COUNT 1
1) 1) "1526919030474-0"
2) 1) "action"
2) "login"
3) "user"
4) "ada"
127.0.0.1:6379> XREVRANGE events + - COUNT 1
1) 1) "1526919030474-1"
2) 1) "action"
2) "purchase"
3) "user"
4) "ada"
5) "item"
6) "book"

XREAD COUNT n STREAMS key id อ่าน entry สูงสุด n รายการจาก stream โดยเริ่ม หลังจาก ID ที่กำหนด การส่ง 0 เป็น ID จะอ่านตั้งแต่ต้น

127.0.0.1:6379> XREAD COUNT 10 STREAMS events 0
1) 1) "events"
2) 1) 1) "1526919030474-0"
2) 1) "action"
2) "login"
3) "user"
4) "ada"
2) 1) "1526919030474-1"
2) 1) "action"
2) "purchase"
3) "user"
4) "ada"
5) "item"
6) "book"

หากต้องการรอ entry ที่ยังไม่มีอยู่ ให้เพิ่ม option BLOCK พร้อมค่า timeout เป็นมิลลิวินาที (ค่า 0 หมายถึงรอไปเรื่อย ๆ) และใช้ $ เป็น ID — $ หมายถึง “เฉพาะ entry ที่ใหม่กว่า entry ล่าสุด ณ เวลาที่ออกคำสั่งนี้”:

127.0.0.1:6379> XREAD BLOCK 0 STREAMS events $
(Blocks until a new entry arrives...)

วิธีนี้เปลี่ยน XREAD ให้กลายเป็น long-poll listener เมื่อ client อื่นรัน XADD events * ... client ที่ block อยู่จะได้รับ entry ใหม่ทันที

XADD events * action "login" user "ada"
XADD events * action "purchase" user "ada" item "book"
XLEN events
XRANGE events - +
XREAD COUNT 10 STREAMS events 0
ตัวเลือกBenefitCost
Streams เก็บทุก entry ไว้อ่านซ้ำได้ตลอด รองรับ consumer ที่มาช้าหรือ worker ที่ต้อง restartใช้หน่วยความจำ/ดิสก์มากกว่า Pub/Sub ถ้าไม่ตั้ง MAXLEN ให้ตัด entry เก่าทิ้ง stream จะโตไม่มีที่สิ้นสุด
XADD ด้วย ID อัตโนมัติ (*)ได้ ID ที่เรียงตามเวลาให้ทันที ไม่ต้องจัดการเองID ผูกกับ timestamp ของเครื่อง Redis หากต้องการ ID ที่มีความหมายทางธุรกิจต้องกำหนดเอง
XREAD BLOCK แบบ long-pollได้ event ใหม่ทันทีโดยไม่ต้อง poll ถี่ ๆconnection ค้างรอระหว่าง block ต้องจัดการ timeout และ reconnect เอง
  • ไม่ตั้ง MAXLEN ให้ stream ที่มี throughput สูง — stream ที่ไม่ trim จะสะสม entry ไปเรื่อย ๆ จนกินหน่วยความจำมหาศาล ควรใช้ XADD key MAXLEN ~ 10000 * ... เพื่อจำกัดขนาด
  • ใช้ XRANGE - + แบบไม่จำกัด COUNT บน production stream ขนาดใหญ่ — ดึงข้อมูลทั้งหมดในครั้งเดียวอาจบล็อก server และใช้ memory มาก ควรระบุ COUNT และ paginate ด้วย ID ที่ได้จาก entry สุดท้าย
  • เข้าใจผิดว่า Streams รับประกัน exactly-once — Redis Streams ให้ at-least-once delivery เท่านั้น หาก consumer ประมวลผลซ้ำได้ (idempotent) ต้องออกแบบ logic ฝั่ง consumer เอง ไม่ใช่พึ่งพา stream ล้วน ๆ

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

Order-event pipeline — ระบบ e-commerce ใช้ Streams เก็บ event อย่าง order_created, payment_confirmed ให้ service หลายตัวอ่านและ replay ได้เมื่อต้อง rebuild state

Activity feed / audit log — แอปที่ต้องเก็บประวัติกิจกรรมผู้ใช้แบบเรียงเวลา ใช้ XADD บันทึกทุก action แล้วให้ dashboard อ่านผ่าน XRANGE ย้อนดูช่วงเวลาที่ต้องการได้

การส่ง `*` เป็น ID ให้กับ XADD ทำอะไร?
argument ใดของ XRANGE ที่หมายถึง 'entry แรกสุดใน stream'?
sequence number ใน entry ID เช่น `1526919030474-3` หมายถึงอะไร?
ความแตกต่างหลักระหว่าง Redis Streams และ Redis Pub/Sub คืออะไร?