สารบัญชุดบทความ

  1. Communication Styles และ Trade-offs
  2. แยก Background Job ก่อน Scale ข้าม Process
  3. ส่ง Job ข้าม Process ด้วย NATS
  4. ทำให้ Job ไม่หายด้วย NATS JetStream
  5. เมื่อ Job เติบโตเป็น Event Stream ด้วย Kafka
  6. เลือก Communication Style ให้เหมาะกับแต่ละ Service

ไม่มี communication style เดียวที่เหมาะกับทั้ง system

ตลอดบทความนี้เราเปลี่ยน solution ตาม business problem ที่เจอ เริ่มจาก function call ใน process เดียว จากนั้นแยก background job ด้วย local queue ส่ง job ข้าม process ด้วย Core NATS เพิ่ม durability ด้วย JetStream และใช้ Kafka เมื่อ event ต้องถูกเก็บและอ่านโดยหลายระบบใน scale ที่ใหญ่ขึ้น

สิ่งสำคัญคือเราไม่ได้กำลังหา broker ที่ดีที่สุดสำหรับทั้ง system เรากำลังเลือกวิธี communicate ให้เหมาะกับ workflow และ ownership ของแต่ละ service

นี่สอดคล้องกับแนวคิดใน Building Microservices ที่มอง communication เป็น design choice แต่ละแบบมี trade-off ของ latency, coupling, reliability และ operational complexity

เริ่มจากคำถามทาง business ก่อนเลือก technology

ก่อนเลือก protocol หรือ broker ให้ถามว่า service นี้กำลังทำอะไร

เรากำลังต้องการคำตอบกลับทันทีหรือไม่

ถ้าไม่มีคำตอบตอนนี้ ระบบสามารถทำงานต่อแล้วแจ้งผลภายหลังได้หรือไม่

มี service เดียวที่ต้องทำ job นี้ หรือหลาย service ต้องรับรู้เหตุการณ์เดียวกัน

ถ้า process ล้มเหลว message หายได้หรือไม่

เราต้อง replay event ย้อนหลังหรือไม่

คำตอบเหล่านี้จะพาเราไปสู่ communication style ที่เหมาะสม

1. Synchronous request สำหรับข้อมูลที่ต้องใช้ทันที

ใช้เมื่อ caller ต้องการ response ก่อนจบ request ตัวอย่างเช่น Order Service ขอ delivery fee จาก Pricing Service เพื่อแสดงยอดให้ลูกค้าเห็นก่อนยืนยัน order

สีใน diagram ใช้เหมือนกันตลอดชุดบทความ: Order / Order API = สีฟ้า; Pricing / Delivery Fee = สีส้ม; Rider Matching = สีเขียว; Payment = สีชมพู; Notification = สีคราม; Restaurant = สีเขียวมะนาว; Analytics = สีม่วง; Data Platform = สีฟ้าอมเขียว; Accounting = สีน้ำตาลอ่อน; Queue / broker / event stream = สีเหลือง; Client / database / file / infrastructure = สีเทา สีบอก service หรือประเภทขององค์ประกอบ โดยชื่อในกล่องระบุหน้าที่ ส่วนกรอบหนาพร้อมคำว่า “เพิ่ม” บอก instance ที่เพิ่มเพื่อ scale ส่วนกรอบแดงเส้นประพร้อมคำว่า “stopped” บอก worker ที่หยุดทำงาน

sequenceDiagram
    box rgb(241,245,249)
        participant App as Customer app
    end
    box rgb(219,234,254)
        participant Order as Order Service
    end
    box rgb(255,237,213)
        participant Pricing as Pricing Service
    end

    App->>Order: Confirm order
    Order->>Pricing: Request delivery fee
    Pricing-->>Order: Return fee
    Order-->>App: Show total price

ข้อดีคือ flow อ่านง่ายและ caller รู้ผลทันที ข้อเสียคือ Order Service ผูกกับ availability และ latency ของ Pricing Service ถ้า Pricing Service ช้า request ของลูกค้าก็ช้าตาม

ไม่ควรใช้ synchronous request กับงานที่ใช้เวลานานหรือไม่จำเป็นต้องตอบใน request เดียว เช่น การหา rider, การส่ง notification หรือการสร้าง daily report

เมื่อ request volume เพิ่มขึ้น เรา scale Order Service และ Pricing Service แยกกันได้ แต่ทุก request ยังต้องผ่าน Pricing Service และ Pricing database อยู่

flowchart LR
    App["Customer app"] --> O1["Order Service 1"]
    App --> O2["Order Service 2"]
    O1 --> P1["Pricing Service 1"]
    O2 --> P2["Pricing Service 2"]
    P1 --> DB["Pricing database"]
    P2 --> DB

    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class App,DB neutral
    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class O1,O2 order
    classDef pricing fill:#ffedd5,stroke:#c2410c,color:#0f172a
    class P1,P2 pricing

ถ้า database เป็นคอขวด การเพิ่ม service instance จะไม่ช่วยเพิ่ม throughput ต่อไป

2. Asynchronous job สำหรับงานที่ทำภายหลังได้

ใช้เมื่อ caller เพียงต้องการฝาก job แล้วให้ worker ทำต่อ ตัวอย่างเช่น Order Service ส่ง FindRider ให้ Rider Matching Service

flowchart LR
    API["Order Service"] --> Broker["Job broker"]
    Broker --> Worker["Rider Matching worker"]
    Worker --> DB["Order database<br/>status = matched"]

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class API order
    classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
    class Broker broker
    classDef rider fill:#dcfce7,stroke:#15803d,color:#0f172a
    class Worker rider
    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class DB neutral

Caller ไม่ต้องรอให้หา rider เสร็จใน request เดียว ระบบสามารถตอบว่า order อยู่ในสถานะ matching จากนั้น worker update status เมื่อทำ job เสร็จ

ใช้ Core NATS เมื่อ message หายได้และต้องการ low latency ใช้ JetStream เมื่อ job ต้องอยู่รอดเมื่อ worker offline และต้องรองรับ retry

Asynchronous job ยังมี trade-off caller ไม่ได้ผลสำเร็จจาก return value ของ function เราจึงต้องมี status, correlation ID, retry policy และวิธีจัดการ duplicate processing

เมื่อ job สะสม เรา scale worker แยกจาก Order Service ได้ ทุก worker ใช้ queue group หรือ durable consumer เดียวกันเพื่อแบ่ง job กันทำ

flowchart LR
    API1["Order API 1"] --> JS["JetStream<br/>FindRider jobs"]
    API2["Order API 2"] --> JS
    JS --> W1["Rider worker 1"]
    JS --> W2["Rider worker 2"]
    JS --> W3["Rider worker 3"]
    W1 --> DB["Order database"]
    W2 --> DB
    W3 --> DB

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class API1,API2 order
    classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
    class JS broker
    classDef rider fill:#dcfce7,stroke:#15803d,color:#0f172a
    class W1,W2,W3 rider
    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class DB neutral

การเพิ่ม worker จะช่วยได้เมื่อ database และ map service ยังมี capacity เหลือ

3. Event notification สำหรับการแจ้งข้อเท็จจริง

ใช้เมื่อ service หนึ่งต้องประกาศว่าเหตุการณ์เกิดขึ้นแล้ว แต่ไม่ควรต้องรู้ว่ามี service ใดนำข้อมูลไปใช้บ้าง

หลังจาก Payment Service ยืนยันการชำระเงิน อาจ publish PaymentCompleted Notification Service ส่งข้อความให้ลูกค้า Restaurant Service เริ่มเตรียมอาหาร Analytics Service เก็บ event สำหรับรายงาน

flowchart LR
    Payment["Payment Service"] --> Event["PaymentCompleted event"]
    Event --> Notification["Notification Service"]
    Event --> Restaurant["Restaurant Service"]
    Event --> Analytics["Analytics Service"]

    classDef payment fill:#fce7f3,stroke:#be185d,color:#0f172a
    class Payment payment
    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class Event neutral
    classDef notification fill:#e0e7ff,stroke:#4338ca,color:#0f172a
    class Notification notification
    classDef restaurant fill:#ecfccb,stroke:#4d7c0f,color:#0f172a
    class Restaurant restaurant
    classDef analytics fill:#ede9fe,stroke:#7c3aed,color:#0f172a
    class Analytics analytics

Event notification ลด direct coupling ระหว่าง publisher กับ subscriber แต่ต้องออกแบบ schema, versioning, ownership และการส่งซ้ำให้ดี

ถ้า event มีเฉพาะ consumer ภายในไม่กี่ตัวและต้องการ latency ต่ำ JetStream อาจเพียงพอ ถ้า event ต้องถูกเก็บนาน, replay บ่อย, มี consumer จำนวนมาก หรือไหลเข้าสู่ data platform Kafka จะเหมาะกว่า

แต่ละ consumer group scale แยกกันได้ตาม processing load ของตัวเอง การเพิ่ม Analytics worker จึงไม่จำเป็นต้องเพิ่ม Notification worker ตาม

flowchart LR
    Event["PaymentCompleted event"] --> N1["Notification consumer 1"]
    Event --> N2["Notification consumer 2"]
    Event --> A1["Analytics consumer 1"]
    Event --> A2["Analytics consumer 2"]
    Event --> A3["Analytics consumer 3"]

    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class Event neutral
    classDef notification fill:#e0e7ff,stroke:#4338ca,color:#0f172a
    class N1,N2 notification
    classDef analytics fill:#ede9fe,stroke:#7c3aed,color:#0f172a
    class A1,A2,A3 analytics

การ scale แบบนี้ต้องแยก consumer group ให้ถูกต้อง ถ้าใช้ group เดียวกันทุก service จะผลัดกันรับ event แทนที่จะได้รับ event ครบทุก service

4. Common data สำหรับการอ่านข้อมูลร่วมกัน

บาง workflow ไม่จำเป็นต้องส่ง request หรือ event ทุกครั้ง หลาย service อาจอ่านข้อมูลจาก data store ที่ตกลง ownership และ schema ร่วมกัน

ตัวอย่างเช่น Accounting Service อ่านไฟล์สรุปยอดการขายที่ Order Service export ตอนสิ้นวัน รูปแบบนี้เหมาะกับ batch workflow ที่ไม่ต้องการ real-time response

flowchart LR
    Order["Order Service"] --> File["Daily settlement file"]
    File --> Accounting["Accounting Service"]

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class Order order
    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class File neutral
    classDef accounting fill:#efdfd0,stroke:#92400e,color:#0f172a
    class Accounting accounting

ข้อดีคือแยกเวลาและ capacity ของแต่ละ service ได้ ข้อเสียคือข้อมูลไม่ real-time และต้องดูแล file format, delivery และการนำไฟล์กลับมาประมวลผลเมื่อ fail

ไม่ควรใช้ common database เป็นทางลัดให้ทุก service แก้ข้อมูลชุดเดียวกัน เพราะ ownership จะไม่ชัดและ schema change ของ service หนึ่งจะกระทบอีก service โดยตรง

เปรียบเทียบ communication style ในภาพรวม

Style ใช้เมื่อ ตัวอย่าง ความเสี่ยงหลัก
Synchronous request ต้องการ response ทันที ขอ delivery fee latency และ dependency สูง
Asynchronous job ทำงานภายหลังได้ หา rider ต้องติดตาม status และ retry
Event notification หลาย service ต้องรู้ว่าเกิดอะไรขึ้น PaymentCompleted schema และ duplicate event
Common data batch หรือ shared reporting workflow daily settlement file ไม่ real-time และต้องดูแล format
Event streaming event volume สูงและต้อง replay order analytics operational complexity สูง

Service เดียวกันไม่จำเป็นต้องใช้ broker เดียวกัน

สมมติว่า food delivery system มี service เหล่านี้

Service Workflow Communication ที่เหมาะสม
Order Service ขอ delivery fee ก่อนยืนยัน order Synchronous request
Rider Matching Service รับ job หา rider JetStream durable consumer
Notification Service รับ OrderStatusChanged JetStream หรือ Core NATS ตามความสำคัญของ event
Payment Service รับ payment request และตอบผล Synchronous request หรือ durable job ตาม payment provider
Analytics Service อ่าน order event จำนวนมากและ replay ย้อนหลัง Kafka consumer group
Accounting Service ประมวลผล settlement ตอนสิ้นวัน Common data หรือ batch event

ทั้งหมดนี้อยู่ใน system เดียวกันได้ แต่ละ service ยังเป็น service แยกกัน มี ownership ของตัวเอง และเลือก communication style ตาม business workflow ของตัวเอง

กรณีที่ไม่ควรใช้ Kafka

Kafka ไม่ควรเป็น default broker สำหรับทุกงาน

ถ้าเป็น request ที่ต้องตอบทันที Kafka ไม่ได้ทำให้ synchronous dependency หายไป ถ้าเป็น job เล็ก ๆ ภายใน service ที่มี worker ไม่กี่ตัว Kafka อาจเพิ่ม deployment, storage และ monitoring โดยไม่จำเป็น ถ้า message มีอายุสั้นและยอมให้หายได้ Core NATS อาจเหมาะกว่า ถ้าต้องการ durable job queue ขนาดไม่ใหญ่ JetStream อาจง่ายกว่าและมี operational cost ต่ำกว่า

ตัวอย่างเช่น FindRider ภายใน food delivery app ไม่จำเป็นต้องใช้ Kafka เพียงเพราะมีจำนวน order สูงขึ้น ถ้า JetStream รองรับ throughput, retry และ retention ที่เราต้องการอยู่แล้ว การใช้ solution ที่เล็กกว่าจะช่วยลดความซับซ้อนของ system

กรณีที่ Kafka เหมาะสม

Kafka เหมาะเมื่อ event กลายเป็น shared data stream ของหลายระบบ

ตัวอย่างเช่น OrderCreated ต้องถูกอ่านโดย Payment, Restaurant, Analytics และ Data Platform แต่ละ consumer group ต้องอ่านด้วยความเร็วของตัวเอง Analytics อาจหยุดไปหลายชั่วโมงแล้ว replay event ที่พลาด ทีม Data อาจอ่าน event ย้อนหลังหลายเดือนเพื่อสร้าง report ใหม่ และ volume ของ event อาจสูงจนต้อง scale ด้วย partition

ในกรณีนี้ Kafka ช่วยให้ event มี retention, replay และ parallel processing ที่ชัดเจน แต่เราต้องยอมรับ operational cost ของ cluster, partition, replication, monitoring และ schema management

ภาพใหญ่ของ architecture

architecture ที่ดีจึงอาจมี communication หลายแบบอยู่ร่วมกัน

flowchart LR
    App["Customer app"] --> Order["Order Service"]
    Order -->|request fee| Pricing["Pricing Service"]
    Order -->|durable FindRider job| JS["NATS JetStream"]
    JS --> Rider["Rider Matching Service"]
    Order -->|OrderCreated event| Kafka["Kafka<br/>order.events"]
    Kafka --> Analytics["Analytics Service"]
    Kafka --> Data["Data Platform"]
    Order -->|daily export| Accounting["Accounting Service"]

    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class App neutral
    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class Order order
    classDef pricing fill:#ffedd5,stroke:#c2410c,color:#0f172a
    class Pricing pricing
    classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
    class JS,Kafka broker
    classDef rider fill:#dcfce7,stroke:#15803d,color:#0f172a
    class Rider rider
    classDef analytics fill:#ede9fe,stroke:#7c3aed,color:#0f172a
    class Analytics analytics
    classDef data fill:#cffafe,stroke:#0e7490,color:#0f172a
    class Data data
    classDef accounting fill:#efdfd0,stroke:#92400e,color:#0f172a
    class Accounting accounting

Order Service ใช้ synchronous request กับ Pricing Service เพราะต้องแสดงราคาให้ลูกค้าทันที ใช้ JetStream กับ Rider Matching Service เพราะเป็น durable job ที่ต้องทำให้สำเร็จ ใช้ Kafka กับ Analytics และ Data Platform เพราะ event ต้องถูกเก็บ, replay และอ่านโดยหลาย consumer group ใช้ daily export กับ Accounting Service เพราะ business requirement เป็น batch ตอนสิ้นวัน

นี่คือเหตุผลที่การออกแบบ microservice ไม่ควรเริ่มจากคำถามว่า “ทีมเราใช้ Kafka หรือ NATS” ควรเริ่มจากคำถามว่า “workflow นี้ต้องการ guarantee และเวลาในการตอบแบบไหน”

เลือก scale service ให้ตรงกับ bottleneck

ลองนึกถึงช่วงเที่ยงที่ order เข้ามาพร้อมกันจำนวนมาก Order Service อาจยังรับ request ได้ แต่ Pricing คำนวณราคาไม่ทัน, job หา rider เริ่มสะสม หรือ Analytics ทำรายงานตาม event ไม่ทัน อาการเหล่านี้ต้องแก้คนละจุด การเพิ่ม instance ให้ทุก service เท่ากันจึงอาจเพิ่ม cost โดยไม่ได้ลดเวลารอของลูกค้า

diagram ในแต่ละหัวข้อแสดงการ scale เฉพาะกรณี โดยสมมติว่าแต่ละ service เริ่มจาก 1 instance จำนวน instance ใช้เพื่ออธิบายแนวคิด ไม่ใช่ sizing สำหรับ production ทุกภาพใช้ สีประจำ service เหมือนกับ diagram อื่นในชุดบทความ instance ใหม่ใช้ กรอบหนาพร้อมคำว่า “เพิ่ม” และ กล่องครอบเส้นประ แสดงขอบเขต service ที่ scale instance เดิมยังใช้สีเดียวกัน เพื่อให้เห็นว่าเป็น service เดียวกัน แต่เพิ่ม capacity เฉพาะส่วนที่ต้องการ

Pricing ช้า: scale ส่วนที่อยู่ในเส้นทางรอคำตอบ

ถ้า latency ของการยืนยัน order สูงขึ้น ให้แยกดูว่าเวลาหายไปที่ Order, Pricing หรือ database ถ้า Pricing ใช้ CPU เต็มจากการคำนวณราคา แต่ database ยังตอบเร็ว การเพิ่ม Pricing instance หลัง load balancer ช่วยแบ่งงานได้ Order Service ไม่จำเป็นต้องเพิ่มตาม หากยังมี capacity รับ request เพียงพอ

flowchart TB
    Order["Order Service: 1 instance"] -->|request fee| LB["Pricing load balancer"]
    subgraph Scale["Pricing: 1 → 2 instances"]
        P1["Pricing 1: เดิม"]
        P2["Pricing 2: เพิ่ม"]
    end
    LB --> P1
    LB --> P2
    P1 --> DB["Pricing database เดิม<br/>ต้องยังมี capacity เหลือ"]
    P2 --> DB

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class Order order
    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class LB,DB neutral
    classDef pricing fill:#ffedd5,stroke:#c2410c,color:#0f172a
    class P1,P2 pricing
    style P2 stroke-width:4px
    style Scale fill:#ffedd5,stroke:#c2410c,color:#0f172a,stroke-width:2px,stroke-dasharray:5 5

แต่ถ้า Pricing รอ query เดียวกันที่ database การเพิ่ม Pricing instance อาจยิ่งเพิ่ม connection และงานที่แย่ง resource กัน ต้องแก้ query, index หรือ capacity ของ database ก่อน และใช้ cache ได้เฉพาะข้อมูลที่ business ยอมรับความเก่าได้ เช่น ราคาที่ขึ้นกับ demand ปัจจุบันอาจใช้ cache นานเท่าข้อมูลพื้นที่จัดส่งไม่ได้

Synchronous request จึง scale แยก service ได้ แต่เวลาตอบลูกค้ายังขึ้นกับ dependency ตลอดเส้นทาง เราต้องกำหนด timeout และจำกัด request ที่กำลังทำงาน เพื่อไม่ให้ความช้าของ Pricing ทำให้ Order ใช้ resource จนหมด การเปลี่ยนเป็น queue อย่างเดียวไม่แก้โจทย์ หากลูกค้ายังต้องเห็นราคาที่ถูกต้องก่อนกดยืนยัน

Rider Matching งานค้าง: scale worker ตามเวลาที่ลูกค้ารอได้

ถ้า Order ตอบรับได้เร็ว แต่เวลารอหา rider เพิ่มขึ้น ให้ดูทั้งจำนวน job ค้างและอายุของ job ที่เก่าที่สุด ในกรณีที่งานคำนวณ matching เป็นคอขวด เราเพิ่ม Rider worker โดยไม่ต้องเพิ่ม Order หรือ Pricing worker หลายตัวแบ่งงานผ่าน durable pull consumer เดียวกันได้ ตามรูปแบบ shared pull consumer ของ JetStream และยังต้องรองรับการประมวลผลซ้ำเมื่อมี retry

flowchart TB
    Order["Order Service: 1 instance"] -->|FindRider| JS["JetStream<br/>shared durable pull consumer"]
    subgraph Scale["Rider Matching: 1 → 3 workers"]
        R1["Rider worker 1: เดิม"]
        R2["Rider worker 2: เพิ่ม"]
        R3["Rider worker 3: เพิ่ม"]
    end
    JS --> R1
    JS --> R2
    JS --> R3
    R1 --> Downstream["Map API / matching database เดิม<br/>จำกัด concurrency ตาม capacity"]
    R2 --> Downstream
    R3 --> Downstream

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class Order order
    classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
    class JS broker
    classDef rider fill:#dcfce7,stroke:#15803d,color:#0f172a
    class R1,R2,R3 rider
    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class Downstream neutral
    style R2 stroke-width:4px
    style R3 stroke-width:4px
    style Scale fill:#dcfce7,stroke:#15803d,color:#0f172a,stroke-width:2px,stroke-dasharray:5 5

สมมติมี 120 jobs ต่อวินาที แต่ worker ทั้งหมดทำได้เพียง 80 jobs ต่อวินาที backlog จะเพิ่มประมาณ 40 jobs ทุกวินาที การมี queue ช่วยพักงานช่วง burst แต่ไม่ได้เพิ่มกำลังประมวลผล ถ้า traffic ระดับนี้อยู่นาน ต้องเพิ่ม capacity ให้ทำได้มากกว่า 120 jobs ต่อวินาที เพื่อรองรับ job ใหม่ที่เข้ามา พร้อมกับลด backlog ที่สะสมอยู่

อย่างไรก็ตาม ถ้าคอขวดคือ rate limit ของ Map API การเพิ่ม worker อาจทำให้ถูกปฏิเสธมากขึ้น ต้องจำกัด concurrency และจัด retry ให้เหมาะสมก่อน และถ้าปัญหาคือไม่มี rider ว่างในพื้นที่ การเพิ่ม compute ก็ไม่ทำให้มี rider เพิ่มขึ้น business ต้องกำหนดเวลารอสูงสุด การแจ้งสถานะ หรือเงื่อนไขยกเลิก order ไว้ด้วย

Asynchronous style ทำให้ scale ฝั่งรับ order กับฝั่งทำงานแยกกันได้ง่ายขึ้น แต่ความสำเร็จต้องวัดจากเวลาที่หา rider เสร็จ ไม่ใช่แค่ API ตอบรับได้เร็ว

Analytics ตามไม่ทัน: scale เฉพาะ consumer group

ถ้า Analytics มี consumer lag สูงขึ้น แต่ Data Platform ยังตาม event ทัน เราสามารถเพิ่มเฉพาะ Analytics consumer โดยไม่ต้องเพิ่ม consumer ของ Data Platform เพราะทั้งสอง group มีตำแหน่งการอ่านของตัวเอง

flowchart TB
    Order["Order Service: 1 instance"] --> Kafka["Kafka: order.events<br/>3 partitions เดิม"]
    subgraph Scale["Analytics group: 1 → 3 consumers"]
        A1["Analytics 1: เดิม"]
        A2["Analytics 2: เพิ่ม"]
        A3["Analytics 3: เพิ่ม"]
    end
    Kafka -->|partition 0| A1
    Kafka -->|partition 1| A2
    Kafka -->|partition 2| A3
    Kafka -->|ทั้ง 3 partitions| Data["Data Platform group<br/>1 consumer เดิม ยังตามทัน"]
    A1 --> Warehouse["Data warehouse เดิม<br/>ต้องยังรับการเขียนทัน"]
    A2 --> Warehouse
    A3 --> Warehouse

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class Order order
    classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
    class Kafka broker
    classDef analytics fill:#ede9fe,stroke:#7c3aed,color:#0f172a
    class A1,A2,A3 analytics
    classDef data fill:#cffafe,stroke:#0e7490,color:#0f172a
    class Data data
    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class Warehouse neutral
    style A2 stroke-width:4px
    style A3 stroke-width:4px
    style Scale fill:#ede9fe,stroke:#7c3aed,color:#0f172a,stroke-width:2px,stroke-dasharray:5 5

สำหรับ Kafka consumer group แบบแบ่ง partition แต่ละ partition ถูกอ่านโดย consumer เดียวใน group ณ เวลาหนึ่ง ถ้า topic มี 3 partitions การเพิ่มจาก 1 เป็น 3 consumers ช่วยแบ่ง partition ได้ แต่ตัวที่ 4 จะไม่มี partition ให้รับผิดชอบในตัวอย่างนี้ ดูรายละเอียดใน Kafka Design

ก่อนเพิ่ม partition ต้องดูการกระจาย key และเงื่อนไข ordering ด้วย ถ้างานหนักกระจุกอยู่ที่ key เดียว การเพิ่ม consumer ไม่ได้แบ่งงานของ key นั้นออกโดยอัตโนมัติ และถ้า Analytics ตันที่การเขียน data warehouse ก็ต้องแก้ปลายทางนั้นก่อน

Event streaming เหมาะกับ business ที่ยอมให้รายงานตามหลังการสั่งอาหารได้ และต้องการ replay จึงแยก capacity ของงานวิเคราะห์ออกจากงานรับ order ได้มากขึ้น แต่ยังต้องเผื่อ broker throughput, storage และ retention ให้รองรับ lag เพราะทุก group ใช้ infrastructure ร่วมกัน

Accounting ปิดยอดไม่ทัน: scale ตาม deadline ของ batch

Accounting ไม่จำเป็นต้องเพิ่ม worker ตามยอด request ช่วงเที่ยงทันที ถ้า business ต้องการผลตอนสิ้นวัน เราสามารถจัดตารางและเพิ่ม batch worker ตามปริมาณไฟล์กับเวลาที่เหลือก่อน deadline เราสามารถแบ่งไฟล์หรือข้อมูลเป็น batch ย่อย แล้วให้ worker ประมวลผลแบบ parallel ได้ ถ้าแต่ละ batch ไม่ต้องรอผลจาก batch อื่น หลังประมวลผลต้องตรวจสอบว่ายอดรวมตรงกับข้อมูลต้นทาง และป้องกันไม่ให้การ retry บันทึกรายการเดิมซ้ำ

flowchart TB
    Order["Order Service: 1 instance"] -->|daily export| Files["Settlement files<br/>แบ่งเป็น batch A และ B ที่ไม่ต้องรอผลกัน"]
    subgraph Scale["Accounting: 1 → 2 workers"]
        B1["Batch worker 1: เดิม"]
        B2["Batch worker 2: เพิ่ม"]
    end
    Files -->|ชุด A| B1
    Files -->|ชุด B| B2
    B1 --> Result["ตรวจยอดรวมและความครบถ้วน<br/>ก่อนปิดยอดตาม deadline"]
    B2 --> Result

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class Order order
    classDef neutral fill:#f1f5f9,stroke:#64748b,color:#0f172a
    class Files neutral
    classDef accounting fill:#efdfd0,stroke:#92400e,color:#0f172a
    class B1,B2,Result accounting
    style B2 stroke-width:4px
    style Scale fill:#efdfd0,stroke:#92400e,color:#0f172a,stroke-width:2px,stroke-dasharray:5 5

ถ้าวันหนึ่ง business ต้องการยอด settlement แบบ real-time เราจะรอรวบรวมข้อมูลไปประมวลผลตอนสิ้นวันไม่ได้อีก การเพิ่ม batch worker อย่างเดียวอาจไม่พอ ต้องทบทวน flow การส่งข้อมูลด้วย

Business character กำหนดว่า scale ได้ง่ายแค่ไหน

ลักษณะของ business ที่มีผลต่อ architecture ไม่ได้มีแค่จำนวน request แต่รวมถึงเวลาที่รอได้, data freshness หรือข้อมูลต้องอัปเดตเร็วแค่ไหน, รูปแบบ traffic และงานใดทำแยกกันได้

ลักษณะของ business Style ที่สอดคล้อง ผลต่อการ scale
ลูกค้าต้องเห็นราคาก่อนยืนยัน Synchronous request เพิ่ม Pricing แยกได้ แต่ต้องดู latency และ capacity ของ dependency ตลอดทาง
รับ order ก่อน แล้วแจ้งผลหา rider ภายหลังได้ Durable asynchronous job ใช้ queue รับ burst และเพิ่ม worker ตาม backlog ภายในเวลารอที่ยอมรับได้
หลายทีมใช้ event เดียวกันและรายงานช้าได้ Event streaming แต่ละ group scale และอ่านย้อนหลังแยกกัน ภายใต้ข้อจำกัด partition และ broker
ปิดยอดตามรอบเวลา Common data / batch ย้ายเวลาประมวลผลและแบ่งงานตาม deadline ได้ โดยต้องตรวจความครบถ้วนของยอด

ถ้านำ Analytics ไปอยู่ใน synchronous chain ของการสร้าง order ทุกครั้งที่ Analytics ตอบช้า ลูกค้าจะรอยืนยัน order นานขึ้นด้วย ทั้งที่ business flow ของการสร้าง order ไม่จำเป็นต้องรอ response จาก Analytics แต่ถ้าย้ายการคำนวณราคาที่ต้องแสดงทันทีไปเป็น asynchronous job หน้า checkout ต้องแสดงสถานะ pending ระหว่างรอ Pricing คำนวณราคา แล้วอัปเดตราคาให้ลูกค้าเห็นเมื่อได้ผล จึงมีขั้นตอนจัดการสถานะและการอัปเดต UI เพิ่มขึ้น ซึ่งอาจทำให้ checkout ซับซ้อนขึ้นโดยไม่จำเป็น

การเลือก architecture style ให้ตรงกับ business จึงช่วยให้ขอบเขตการ scale ตรงกับงานที่โตจริง งานที่ต้องตอบทันทีได้ capacity ตาม latency ที่ต้องรักษา งานที่รอได้มี buffer และ worker ของตัวเอง ส่วนงานรายงานหรือ batch สามารถเลือกเวลาประมวลผลและปรับจำนวน worker ให้เสร็จทัน deadline โดยควบคุม cost ได้

ระบบจะ scale ได้ง่ายขึ้นเมื่อแต่ละ service เพิ่ม capacity ได้ตาม workload ของตัวเอง และ dependency รองรับ load ที่เพิ่มขึ้นได้ การเพิ่ม broker หรือแยก microservice มากขึ้นเพียงอย่างเดียวจึงไม่ได้ทำให้ scale ง่ายขึ้นเสมอไป

สรุป

Synchronous request เหมาะกับข้อมูลที่ caller ต้องใช้ทันที Asynchronous job เหมาะกับงานที่ทำภายหลังได้ Event notification เหมาะกับการกระจายข้อเท็จจริงไปยังหลาย service Common data เหมาะกับ batch workflow ที่ตกลง format และ ownership ร่วมกัน JetStream เหมาะกับ durable job และ event ภายในระบบที่ต้องการ retry Kafka เหมาะกับ event stream ขนาดใหญ่ที่ต้อง replay, partition และเชื่อมต่อหลายระบบ

ไม่มี style ใดดีที่สุดในทุกสถานการณ์ การเลือกที่ดีคือการยอมรับ trade-off อย่างชัดเจน และเลือกความซับซ้อนเท่าที่ business requirement ต้องการ