บทที่ 9 · Part 2 — Building Blocks

Queues, Streams & Schedulers

Queue ต่างจาก stream อย่างไร, delivery guarantee, DLQ, consumer lag และงาน scheduled ที่ยากกว่าที่คิด

"เอา Kafka มาใส่" เป็นคำตอบที่ได้ยินบ่อยที่สุดและอธิบายน้อยที่สุด Queue กับ stream แก้ปัญหาต่างกัน มี guarantee ต่างกัน และล้มต่างกัน บทนี้แยกให้ชัด แล้วต่อด้วยเรื่องที่ยากกว่าที่คนคิด — งานที่ตั้งเวลาไว้

จบบทนี้คุณจะ

  • แยก queue กับ log/stream ได้ และเลือกถูกตัว
  • พูดเรื่อง delivery guarantee ได้อย่างซื่อสัตย์ รวมถึงความจริงเกี่ยวกับ "exactly-once"
  • ออกแบบ DLQ, poison-message handling และ consumer lag เป็น SLI ได้
  • รู้ว่าทำไม "รันทุก 5 นาที" กับ "รันครั้งเดียวในอนาคตต่อผู้ใช้ 10 ล้านคน" เป็นปัญหาต่างกันโดยสิ้นเชิง

Queue กับ stream ไม่ใช่เครื่องมือเดียวกัน

Queue (SQS, RabbitMQ)Log / stream (Kafka, Kinesis, Pulsar)
Modelงานถูก consume แล้วหายไปappend-only log; consumer ถือ offset ของตัวเอง
Replayไม่ได้ (ack แล้วหายเลย)ได้ — ถอย offset แล้วประมวลผลใหม่
Consumerworker หลายตัวแย่งงานจาก queue เดียวconsumer group อิสระหลายกลุ่มอ่านข้อมูลชุดเดียวกัน
Orderingbest effort (หรือโหมด FIFO ที่ throughput ลด)เข้มงวดต่อ partition
Retentionจนกว่าจะถูก consumeตามเวลาหรือขนาด (ชั่วโมงถึงตลอดไป)
เหมาะกับการกระจายงาน: ส่งอีเมล, สร้าง PDF, เรียก railevent backbone: subscriber หลายราย, audit trail, สร้าง read model ใหม่, analytics

ถามคำถามเดียวก่อนเลือก

"มีใครอีกที่จะต้องอ่านข้อมูลชุดนี้ในอนาคต — และเราจะต้องอ่านย้อนหลังไหม" ถ้าคำตอบคือใช่ทั้ง 2 ข้อ ใช้ log · ถ้าเป็นงานที่ทำเสร็จแล้วจบ ใช้ queue การใช้ log เป็น work queue ทำได้ แต่คุณจะเจอเรื่อง rebalance กับ offset management ที่ queue ไม่มี

Delivery guarantee พูดกันตรงๆ

Guaranteeความจริงใช้ได้เมื่อ
At-most-onceอาจหายข้อมูลที่ทิ้งได้จริงๆ เท่านั้น
At-least-onceอาจซ้ำนี่คือของที่คุณได้จริงบน production — ออกแบบ consumer ให้ idempotent แล้วปลอดภัย
"Exactly-once"จริงภายในขอบเขต transaction ของระบบเดียว (เช่น Kafka read-process-write with transactions) แต่ไม่มีจริงแบบ end-to-end ข้าม network ไปหา third party—

ประโยคที่ควรใช้ในห้องรีวิว

"at-least-once delivery + idempotent processing = effectively-once outcome"

พูดแบบนี้แล้วใครที่เคยรันระบบ payment จะเชื่อคุณมากขึ้น ส่วนคนที่บอกว่าระบบเขา exactly-once end-to-end กับธนาคาร คือคนที่ยังไม่เจอ timeout จริง

เรื่อง operational ที่ต้องออกแบบตั้งแต่ต้น

1. Dead-letter queue + runbook

DLQ ที่ไม่มีใครดูคือเครื่องทำข้อมูลหายแบบเงียบ

ต้องมี 3 อย่างพร้อมกัน:

  • Alert เมื่อ DLQ depth > 0 (ไม่ใช่ > 100)
  • ระบุว่าใครเป็นคนตรวจ และตรวจภายในกี่ชั่วโมง
  • ขั้นตอนการ replay message หลังแก้ bug แล้ว — เขียนไว้ก่อนขึ้น production

สำหรับ message ที่เกี่ยวกับเงิน การ replay ต้องมีคนอนุมัติ ไม่ใช่ปุ่มที่ใครกดก็ได้

2. Poison-message handling

Message ที่ผิดรูป 1 ใบห้ามบล็อก partition ไปตลอดกาล จำกัดจำนวน retry แล้วส่งเข้า DLQ

กฎที่ใช้ได้:
  attempt 1-5    retry แบบ exponential backoff + jitter
  attempt 6      -> DLQ พร้อม metadata: error, stack, offset, timestamp
  commit offset  -> partition เดินต่อได้

--> จุดสำคัญ: commit offset ต่อจากนั้น ไม่ใช่ค้างอยู่ที่ offset เดิม

ช่องว่างระหว่าง "ส่งเข้า DLQ" กับ "commit offset"

ถ้า process ตายหลังส่งเข้า DLQ แต่ก่อน commit offset → message จะถูกประมวลผลซ้ำ (ยอมรับได้ ถ้า DLQ consumer dedupe ด้วย event_id)

ถ้าสลับลำดับเป็น commit offset ก่อนแล้วส่ง DLQ → message หายจริง เมื่อ process ตายระหว่าง 2 ขั้น

ทางเลือกที่ปลอดภัย เรียงตามความน่าใช้:

  1. ส่ง DLQ ก่อน แล้วค่อย commit offset และให้ DLQ dedupe ด้วย event_id (at-least-once ที่ DLQ ดีกว่าข้อมูลหาย)
  2. DLQ เป็นตารางในฐานข้อมูลเดียวกับ state ของ consumer แล้ว insert DLQ row + advance offset ใน transaction เดียว (ใช้ได้เมื่อ consumer เก็บ offset ของตัวเองใน DB)
  3. ใช้ transactional read-process-write ของ broker ถ้ามีให้ใช้ และอยู่ในขอบเขตระบบเดียวกัน

ห้ามสรุปว่า "แค่ส่ง DLQ" คือจบ — ให้เขียนใน design doc ว่าเลือกข้อไหนและเพราะอะไร

3. Consumer lag เป็น SLI ระดับหนึ่ง

วัด lag เป็นเวลา ("เราตามหลังอยู่ 4 นาที") ไม่ใช่แค่จำนวน message เพราะคนฝ่ายธุรกิจเข้าใจหน่วยวินาที ไม่เข้าใจว่า 40,000 message หมายถึงอะไร

4. การเลือก partition / shard key

Key กำหนดทั้ง ordering และ parallelism พร้อมกัน

เลือก key แบบได้เสีย
account_idevent ของ account เดียวกันเรียงลำดับกันaccount ที่ร้อนมากทำให้ partition เดียวตัน
สุ่ม / round robinโหลดกระจายสมบูรณ์ไม่มีการรับประกันลำดับเลย
merchant_idเรียงลำดับต่อ merchantmerchant ใหญ่ 1 ราย = hot partition
hash(account_id) % Nกระจายดีและยังเรียงต่อ accountเปลี่ยน N แล้ว key ย้าย partition — วางแผน N ให้เผื่อไว้ตั้งแต่ต้น

จำนวน partition คือเพดาน parallelism ที่แก้ทีหลังแพง

consumer ใน 1 group มากกว่าจำนวน partition = ตัวที่เกินมานั่งว่าง เพิ่ม partition ทีหลังทำได้ แต่ key จะถูก map ไป partition ใหม่ ซึ่งทำลายการรับประกันลำดับสำหรับ key ที่ย้าย ให้ตั้งจำนวน partition เผื่อการเติบโต 5–10 ปี ตั้งแต่วันแรก (มันถูกกว่าที่คิด)

5. พฤติกรรมตอน backlog

ถ้า consumer ตามหลังไป 2 ชั่วโมง จะเกิดอะไร ต้องตอบ 3 ข้อ:

  • เพดาน storage และ TTL — message หายไหมถ้า retention หมดก่อน consume
  • message ยังมีความหมายไหมตอนที่ได้ประมวลผลจริง — "ห้ามส่ง OTP ที่อายุ 3 ชั่วโมง"
  • จะระบายอย่างไร — เพิ่ม consumer ได้ถึงเพดาน partition, หรือต้องทิ้งบางส่วน (และถ้าทิ้ง ใครอนุมัติ)
Message ที่หมดอายุแล้ว — ต้องมี guard ที่ consumer

if now() - event.occurred_at > max_useful_age {
    metric.expired_dropped.inc()
    log.warn("dropping stale event", event_id, age)
    return ack        // ack ทิ้ง อย่า retry
}

--> ค่า max_useful_age ต่างกันตามชนิด event:
    OTP push        = 2 นาที
    push แจ้งเตือน  = 1 ชั่วโมง
    ledger event    = ไม่มีวันหมดอายุ ต้องประมวลผลให้ได้ทุกใบ

Event เรื่องเงินไม่มีวันหมดอายุ

การ drop event ที่กระทบยอดเงินหรือ ledger เพราะ "มันเก่าแล้ว" คือการทำให้บัญชีไม่ตรง — ต้องประมวลผลให้ครบ แม้จะช้า แยกให้ชัดใน design ว่า topic ไหน drop ได้และ topic ไหน drop ไม่ได้

[!TIP] Real case — DoorDash เปลี่ยนจาก RabbitMQ ไป Kafka ระบบประมวลผล task แบบ asynchronous ของ DoorDash สร้างบน Celery + RabbitMQ ตอนโหลดขึ้นพวกเขาเจอ broker ไม่เสถียร, ความล้มเหลวที่คาดเดาไม่ได้ตอน scale up และ outage ที่ผูกกับการที่ broker เป็นจุดประสานงานจุดเดียว พวกเขาย้ายไป Kafka พร้อมสร้าง "Async Service" ขึ้นมาเฉพาะ และรายงานว่ากำจัด outage คลาสนั้นไปได้หมด

บทเรียน: broker คือ infrastructure และการเลือก infrastructure มีเพดานการ scale ที่โผล่มาในรูปของปัญหา availability ไม่ใช่ปัญหา throughput

และข้อสังเกตสำคัญ: พวกเขาไม่ได้แค่สลับ broker — พวกเขาออกแบบเส้นทางการส่ง task ใหม่ การสลับ component โดยไม่เปลี่ยน design มักแค่ย้ายที่เจ็บ — Eliminating Task Processing Outages by Replacing RabbitMQ with Apache Kafka

งานที่ตั้งเวลาไว้ — ยากกว่าที่คิด

"รันทุก 5 นาที" กับ "รันครั้งเดียวในอนาคตให้ผู้ใช้แต่ละคนจาก 10 ล้านคน" เป็นปัญหาต่างกันโดยสิ้นเชิง

Periodic job

ต้องมี 4 อย่าง:

  1. Leader election หรือ distributed lock เพื่อไม่ให้ N replica รันพร้อมกัน
  2. Jitter เพื่อไม่ให้ทุกอย่างยิงที่ :00
  3. Overrun policy — ถ้ารอบก่อนยังไม่จบ จะ skip หรือ queue ต่อ
  4. Idempotent เพราะมันจะถูก retry แน่นอน
Cron alignment — ปัญหาที่มองไม่เห็นจนกระทั่งมันเกิด

02:00  งาน settlement เริ่ม
02:00  งาน backup เริ่ม
02:00  งาน reindex search เริ่ม
02:00  งานล้าง log เริ่ม
--> DB CPU 100% เป็นเวลา 40 นาที ทุกคืน
    และไม่มีใครรู้เพราะไม่มีใครดู dashboard ตอน 02:00

แก้: กระจายเป็น 02:00 / 02:17 / 02:33 / 03:05
     และใส่ jitter สุ่ม 0-300 s ให้ทุกงาน

Per-entity timer ในระดับใหญ่

1 แถวต่อ 1 timer พร้อม index บน due_at และ worker ที่ poll — วิธีนี้พาคุณไปได้ไกลกว่าที่คิด (ระดับหลายล้าน timer)

-- Worker poll pattern ที่ใช้ได้จริงถึงหลายล้าน timer
BEGIN;
  SELECT timer_id, entity_id, payload, attempts
    FROM timers
   WHERE due_at <= now()
     AND status = 'PENDING'
     AND (lease_expires_at IS NULL OR lease_expires_at < now())  -- reclaim ของที่ค้าง
   ORDER BY due_at
   LIMIT 500
     FOR UPDATE SKIP LOCKED;      -- worker หลายตัวไม่แย่งแถวเดียวกัน

  UPDATE timers
     SET status           = 'CLAIMED',
         claimed_at       = now(),
         lease_expires_at = now() + interval '2 minutes',  -- lease ไม่ใช่ล็อกถาวร
         attempts         = attempts + 1
   WHERE timer_id = ANY($1);
COMMIT;
-- ทำงาน (ต้อง idempotent) แล้วจึง UPDATE status = 'DONE'

FOR UPDATE SKIP LOCKED คือหัวใจ — มันทำให้ worker หลายตัวดึงงานคนละชุด โดยไม่ต้องมี coordinator

`CLAIMED` โดยไม่มี lease = timer ค้างตลอดกาล

ถ้า worker ตายหลังเปลี่ยนสถานะเป็น CLAIMED แต่ก่อนทำงานเสร็จ timer นั้นจะไม่มีใครหยิบอีกเลย และไม่มี error ที่ไหน — งานที่ควรเกิดขึ้นก็แค่ไม่เกิด

ต้องมี 4 อย่างครบ:

  1. lease_expires_at ไม่ใช่แค่ claimed_at — และ query หยิบงานต้อง reclaim แถวที่ lease หมดอายุ (ตามโค้ดข้างบน)
  2. Reaper / reclaim ที่ทำงานได้เอง (ในกรณีนี้ฝังอยู่ใน query แล้ว) บวก metric count(*) WHERE status='CLAIMED' AND lease_expires_at < now()
  3. Retry policy พร้อมเพดาน attempts แล้วส่งเข้า dead-letter table เมื่อเกิน — ไม่ใช่ retry ไม่จำกัด
  4. การทำงานต้อง idempotent เพราะ lease ที่หมดอายุขณะที่ worker เดิม ยังทำงานอยู่ทำให้เกิดการรันซ้อน (ปัญหาเดียวกับ distributed lock ในบทที่ 10 — ถ้างานแตะเงิน ความถูกต้องต้องอยู่ที่ constraint ในฐานข้อมูล ไม่ใช่ที่ lease)

เกินกว่านั้นให้ใช้ partition ตาม time bucket, delay queue, หรือ durable execution engine

Durable workflow

สำหรับ business process หลายขั้นที่ยาว มี retry, timer และ compensation (การขอสินเชื่อ, การตรวจ KYC, การจัดการข้อพิพาท) workflow engine อย่าง Temporal มักดีกว่า state machine ที่เขียนมือใน cron job เพราะได้ retry, timeout และ history มาฟรี — ถ้าสนใจเรื่องนี้ คอร์สTemporal for Fintech ในเว็บนี้ลงรายละเอียดทั้งหมด

นาฬิกาไม่เคยตรงกัน

ห้ามใช้ "timestamp ที่ใหม่กว่าชนะ" เป็นกลไกความถูกต้อง

Wall clock drift และกระโดดได้ (NTP ปรับ, leap second, VM ถูก suspend)

  • ใช้ monotonic clock สำหรับวัดระยะเวลา
  • ห้ามใช้การเทียบ timestamp ข้ามเครื่องเป็นตัวตัดสินว่า write ไหนชนะ — ใช้ version number, sequence จาก DB, หรือ logical clock

Checklist ก่อนขึ้น production

  • เลือก queue หรือ log โดยตอบคำถาม "ต้อง replay ไหม / มีผู้อ่านหลายรายไหม" ได้
  • Consumer idempotent ทุกตัว และมีที่เก็บ event_id ที่ประมวลผลแล้ว
  • DLQ มี alert ที่ depth > 0, มีเจ้าของ, และมีขั้นตอน replay ที่เขียนไว้
  • Retry มีเพดาน มี backoff + jitter และไม่บล็อก partition
  • Consumer lag วัดเป็นเวลาและมี alert
  • จำนวน partition เผื่อการเติบโตแล้ว และ key ที่เลือกไม่มี hot partition ที่รู้ล่วงหน้า
  • มี guard เรื่อง message หมดอายุ และระบุชัดว่า topic ไหน drop ไม่ได้
  • Periodic job มี leader election, jitter, overrun policy และ idempotent
  • ไม่มี cron ตัวไหนตั้งไว้ตรง :00 พร้อมกันหลายตัว
  • ไม่มีที่ไหนใช้ "timestamp ใหม่กว่าชนะ" เป็นกลไกความถูกต้อง

สรุปบทนี้

Queue กระจายงาน log เก็บข้อเท็จจริงให้อ่านซ้ำได้ · ที่ขอบเขตของ side effect ให้ออกแบบเป็น at-least-once เสมอ เว้นแต่พิสูจน์ได้ว่า guarantee และขอบเขตของ transaction ครอบคลุมถึงตรงนั้น ซึ่งหมายความว่าในทางปฏิบัติ consumer ต้อง idempotent · DLQ ที่ไม่มีใครดูคือข้อมูลหายแบบเงียบ และการส่ง DLQ กับ commit offset ต้องคิดเรื่องลำดับให้ชัด · timer ที่ CLAIMED ต้องมี lease + reaper · partition key กำหนดทั้งลำดับและ parallelism · event เรื่องเงิน drop ไม่ได้ · และนาฬิกาไม่เคยตรงกันพอที่จะใช้ตัดสินความถูกต้อง