สารบัญชุดบทความ
- Communication Styles และ Trade-offs
- แยก Background Job ก่อน Scale ข้าม Process
- ส่ง Job ข้าม Process ด้วย NATS
- ทำให้ Job ไม่หายด้วย NATS JetStream
- เมื่อ Job เติบโตเป็น Event Stream ด้วย Kafka
- เลือก Communication Style ให้เหมาะกับแต่ละ Service
Core NATS แก้ปัญหาได้ แต่ job ยังหายได้
ในตอนที่แล้ว Order API ส่ง FindRider ข้าม process ผ่าน Core NATS ได้แล้ว
เราสามารถเพิ่ม Rider Matching worker ได้โดยไม่ต้องให้ API รู้จัก worker แต่ละตัว
แต่ถ้า worker ทุกตัวหยุดอยู่ในช่วงที่ API publish message Core NATS จะไม่มี subscriber รับ message นั้น เมื่อไม่มีใครรับ message ก็ไม่มีที่เก็บไว้เพื่อรอ worker ที่กลับมา online
สีใน diagram ใช้เหมือนกันตลอดชุดบทความ: Order / Order API = สีฟ้า; Rider Matching = สีเขียว; Queue / broker / event stream = สีเหลือง; Client / database / file / infrastructure = สีเทา สีบอก service หรือประเภทขององค์ประกอบ โดยชื่อในกล่องระบุหน้าที่ ส่วนกรอบแดงเส้นประพร้อมคำว่า “stopped” บอก worker ที่หยุดทำงาน
flowchart LR
API["Order API<br/>publish FindRider"] --> NATS["Broker: Core NATS<br/>subject: rider.find"]
subgraph Group["Queue group"]
W1["✕ Matching worker 1<br/>stopped"]
W2["✕ Matching worker 2<br/>stopped"]
end
NATS -.->|message lost| W1
NATS -.->|message lost| W2
classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
class API order
classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
class NATS broker
classDef rider fill:#dcfce7,stroke:#15803d,color:#0f172a
class W1,W2 rider
style W1 stroke:#dc2626,stroke-width:3px,stroke-dasharray:5 5
style W2 stroke:#dc2626,stroke-width:3px,stroke-dasharray:5 5
style Group fill:#dcfce7,stroke:#15803d,color:#0f172a,stroke-width:2px
สำหรับ FindRider นี่ไม่ใช่พฤติกรรมที่เราต้องการ
order ยังต้องมี rider แม้ worker จะ restart หรือมี network failure ชั่วคราว
เราต้องเก็บ message ไว้จนกว่า worker จะพร้อม
เราต้องเพิ่มชั้นที่ทำหน้าที่เก็บ message แบบ durable
เมื่อ API publish FindRider ระบบควร save message ไว้ก่อน
จากนั้น worker จะมาอ่านเมื่อพร้อม
ถ้า worker อ่าน message แล้ว process สำเร็จ ระบบควรรู้ว่า job เสร็จแล้ว ถ้า worker crash ก่อนยืนยันผล ระบบควรส่ง message เดิมกลับมาให้ทำใหม่ได้
flowchart LR
API["Order API"] --> JS["Durable message store"]
JS --> W1["Matching worker 1"]
JS --> W2["Matching worker 2"]
W1 -->|ack after success| JS
W2 -->|ack after success| JS
classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
class API order
classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
class JS broker
classDef rider fill:#dcfce7,stroke:#15803d,color:#0f172a
class W1,W2 rider
แนวคิดนี้ต่างจาก Core NATS ตรงที่ message ไม่ได้ผูกอยู่กับช่วงเวลาที่ subscriber online เท่านั้น
ระบบสามารถเก็บสถานะของ message ว่ายังรอส่ง, ถูกส่งให้ worker แล้ว, หรือได้รับ ack แล้ว
JetStream เพิ่ม durable storage ให้ NATS
NATS JetStream คือความสามารถของ NATS ที่เพิ่มการเก็บ message และติดตามความคืบหน้าของการอ่าน
เราสามารถสร้าง Stream เพื่อบอกว่า message subject ไหนต้องถูกเก็บ
เช่น stream RIDER_JOBS อาจรับทุก message ที่ส่งไปยัง rider.find
flowchart LR
API["Order API<br/>publish rider.find"] --> JS["NATS JetStream<br/>Stream: RIDER_JOBS"]
JS --> C["Durable consumer<br/>rider-matchers"]
C --> W["Rider Matching worker"]
classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
class API order
classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
class JS,C broker
classDef rider fill:#dcfce7,stroke:#15803d,color:#0f172a
class W rider
Stream จึงเป็นพื้นที่เก็บ message ตาม retention policy
ส่วน Consumer คือสถานะการอ่านของ worker ว่าอ่านถึง message ไหนแล้ว
การแยกสองสิ่งนี้ทำให้ message หนึ่งชุดถูกเก็บไว้ใน stream ขณะที่ worker สามารถหยุดและกลับมาอ่านต่อจาก consumer เดิมได้
Worker ต้อง ack หลังทำ job สำเร็จ
ลำดับสำคัญอยู่ที่จังหวะของ ack
worker ควรส่ง ack หลังหา rider และ save status ของ order สำเร็จแล้ว
sequenceDiagram
box rgb(219,234,254)
participant API as Order API
end
box rgb(254,243,199)
participant JS as NATS JetStream
end
box rgb(220,252,231)
participant W as Matching worker
end
box rgb(241,245,249)
participant DB as Order database
end
API->>JS: Publish FindRider
JS-->>API: Publish accepted
Note over JS: Store message durably
W->>JS: Fetch message
W->>W: Find rider
W->>DB: Save rider + order status
DB-->>W: Save succeeded
W->>JS: Ack message
Note over JS: Mark message as processed
Publish accepted หมายถึง JetStream รับและเก็บ message แล้ว
ยังไม่ได้หมายความว่า order จับคู่ rider สำเร็จ
ผลสำเร็จของ business operation เกิดขึ้นเมื่อ worker save status และส่ง ack หลังจากนั้น
ถ้า worker ล้มเหลวก่อนส่ง ack
สมมติว่า worker อ่าน order-3 แล้ว แต่ process crash ก่อน save status หรือก่อนส่ง ack
JetStream จะยังไม่ถือว่า message เสร็จ
เมื่อถึงเวลาที่กำหนด ระบบจึงสามารถ redeliver message ให้ worker ทำใหม่ได้
sequenceDiagram
box rgb(254,243,199)
participant JS as NATS JetStream
end
box rgb(220,252,231)
participant W as Matching worker
end
JS->>W: Deliver order-3
W->>W: Process job
Note over W: Process crashes before ack
JS-->>W: Redeliver order-3
W->>W: Process job again
W->>JS: Ack order-3
พฤติกรรมนี้เรียกว่า at-least-once delivery
message จะถูกส่งอย่างน้อยหนึ่งครั้ง แต่มีโอกาสถูกส่งซ้ำ
ดังนั้น worker ต้องออกแบบให้การ process ซ้ำไม่ทำให้ order เสียสถานะ
วิธีหนึ่งคือใช้ order_id หรือ job_id เป็น idempotency key
ก่อน update database ให้ตรวจว่า job นี้ถูกทำสำเร็จไปแล้วหรือยัง
ถ้าทำไปแล้วก็ไม่ต้องสร้างผลลัพธ์ซ้ำ แต่ยังสามารถส่ง ack เพื่อปิด message ได้
เราได้ durability แต่ต้องดูแล state เพิ่มขึ้น
JetStream แก้ปัญหา message หายเมื่อ worker offline ได้ แต่ระบบมีส่วนที่ต้องออกแบบเพิ่ม
- กำหนด retention และ storage ของ
Stream - กำหนด
ack timeoutและจำนวนครั้งที่ retry - monitor message ที่ค้างและ message ที่ถูก redeliver
- ทำให้ worker รองรับ duplicate processing
- update order status ให้สอดคล้องกับผลของ job
ตอนนี้เรามี durable job queue สำหรับ FindRider แล้ว
ตอนถัดไปเราจะดูกรณีที่หลาย service ต้องอ่าน event เดียวกัน
เพราะ queue group เหมาะกับการแบ่ง job ให้ worker หนึ่งตัวทำ
แต่บาง event ต้องถูกส่งให้ทุก service ที่ subscribe อยู่