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

RabbitMQ Streams

บทเรียนก่อนหน้าวาง RabbitMQ กับ Kafka เป็นขั้วตรงข้าม — queue แบบ delete-on-ack เทียบกับ log ที่เก็บไว้ RabbitMQ Streams ทำให้เส้นแบ่งนั้นเบลอ: Streams นำ log แบบ append-only สไตล์ Kafka เข้ามา ใน RabbitMQ เอง คุณจึงได้ replay และ fan-out สูง ๆ โดยไม่ต้องเพิ่มระบบตัวที่สอง

stream คือ queue type ตัวหนึ่ง (declare ด้วย x-queue-type: stream) แต่ทำงานต่างจาก classic queue มาก:

  • message ถูก append เข้า log และ เก็บไว้ ไม่ถูกลบตอนอ่าน
  • การอ่านเป็นแบบ non-destructive — consume ไม่ลบอะไรทั้งนั้น
  • consumer แต่ละตัวอ่านจาก offset ที่ตัวเองควบคุม และเริ่มจากต้น, ท้าย, timestamp หรือตำแหน่งที่ save ไว้ก็ได้
  • consumer หลายตัวอ่าน stream เดียวกัน ได้อย่างอิสระและด้วยความเร็วต่างกัน
flowchart LR
  prod["Producer"] -->|append| log["Stream log
[0][1][2][3][4][5]"]
  log --> c1["Consumer A
offset 1 (กำลัง replay)"]
  log --> c2["Consumer B
offset 5 (live)"]
  log --> c3["Consumer C
จาก timestamp"]
stream: log เดียวที่เก็บไว้ consumer หลายตัวที่ offset ของตัวเอง
Classic queueStream
ตอนอ่านถูกลบหลัง ackถูกเก็บไว้ (non-destructive)
Replayไม่ได้ได้ — reset offset
ผู้อ่านหลายตัวของ data เดียวกันแข่งกันแย่ง messageแต่ละตัวอ่านทั้ง log อย่างอิสระ
Retentionจนกว่าจะถูก consumeตามเวลาหรือขนาด (เช่น เก็บ 7 วัน)
เก่งที่สุดตรงtask distribution, งานระดับ per-messageevent fan-out, replay, การอ่านระดับใหญ่

เพราะ stream เก็บ data ไว้ คุณจึง declare retention — ตามเวลาหรือขนาดรวม — แทนที่จะพึ่งการ consume มาเคลียร์ segment เก่าจะถูกตัดทิ้งเมื่อเกินอายุ

เลือก stream เมื่อคุณอยากได้ log สไตล์ Kafka แต่ไม่อยากรัน Kafka: คุณต้องการ replay, คุณมี consumer หลายตัวอ่าน event เดียวกัน, หรือคุณกำลังรับ firehose ที่ throughput สูง และอยากได้ retention ตามเวลา และยังมี super stream — stream ที่ partition แล้ว — สำหรับ scale stream เชิง logic ตัวเดียวข้าม node คล้าย partition ของ Kafka

อยู่กับ classic (หรือ quorum) queue เมื่อ message คือ task ที่ทำครั้งเดียวแล้วทิ้ง, เมื่อคุณต้องการ routing และ retry ระดับ per-message หรือเมื่อไม่ต้องการประวัติ งาน task-distribution และ RPC ส่วนใหญ่ยังต้องการ queue ไม่ใช่ stream

เกิดอะไรขึ้นกับ message ใน RabbitMQ stream เมื่อ consumer อ่าน?
consumer เลือกจุดเริ่มอ่าน stream อย่างไร?
data ถูกเคลียร์ออกจาก stream อย่างไร?
workload ไหนที่ยังควรอยู่บน classic/quorum queue มากกว่า stream?