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

  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

JetStream แก้ปัญหา job หายแล้ว แต่ระบบเริ่มมีคำถามใหม่

ในตอนที่แล้วเราใช้ NATS JetStream เก็บ FindRider แบบ durable ถ้า worker offline message ยังอยู่ใน stream และ worker กลับมาอ่านต่อได้

สำหรับ background job ภายใน service นี่เป็น solution ที่เหมาะสม แต่เมื่อ food delivery app โตขึ้น message ที่เราต้องส่งเริ่มไม่ได้มีแค่ job สำหรับ worker ตัวเดียว

หลังจากลูกค้าสร้าง order ระบบอาจต้องแจ้งหลาย service

  • Payment Service ตรวจสอบการชำระเงิน
  • Restaurant Service อัปเดตสถานะการเตรียมอาหาร
  • Rider Matching Service หา rider
  • Notification Service ส่ง push notification
  • Analytics Service เก็บข้อมูลเพื่อวิเคราะห์

ทุก service ต้องการข้อมูลจากเหตุการณ์เดียวกัน แต่แต่ละ service ต้องอ่านด้วยความเร็วและช่วงเวลาที่ต่างกัน

จากการแบ่ง job สู่การกระจาย event

FindRider เป็น job ที่ต้องการ worker เพียงหนึ่งตัวทำงาน เราใช้ queue group เพื่อแบ่ง job ระหว่าง worker

แต่ OrderCreated เป็น event ที่หลาย service ต้องได้รับ โดยไม่จำเป็นต้องอ่านในเวลาเดียวกัน Payment Service ไม่ควรแย่ง event เดียวกับ Analytics Service แต่ละ service ควรมี consumer หรือโปรแกรมที่อ่าน event ไปประมวลผลของตัวเอง ส่วน Order Service ที่เขียน event เข้ามา เราเรียกว่า producer

สีใน diagram ใช้เหมือนกันตลอดชุดบทความ: Order / Order API = สีฟ้า; Payment = สีชมพู; Restaurant = สีเขียวมะนาว; Analytics = สีม่วง; Queue / broker / event stream = สีเหลือง สีบอก service หรือประเภทขององค์ประกอบ โดยชื่อในกล่องระบุหน้าที่

flowchart LR
    Order["Order Service<br/>publish OrderCreated"] --> Stream["Event stream"]
    Stream --> Payment["Payment Service<br/>consumer"]
    Stream --> Restaurant["Restaurant Service<br/>consumer"]
    Stream --> Analytics["Analytics Service<br/>consumer"]

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class Order order
    classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
    class Stream broker
    classDef payment fill:#fce7f3,stroke:#be185d,color:#0f172a
    class Payment payment
    classDef restaurant fill:#ecfccb,stroke:#4d7c0f,color:#0f172a
    class Restaurant restaurant
    classDef analytics fill:#ede9fe,stroke:#7c3aed,color:#0f172a
    class Analytics analytics

นี่คือจุดที่ problem เปลี่ยนจาก ส่ง job ให้ worker ทำ เป็น เก็บ event เพื่อให้หลาย service นำไปใช้ เราต้องการ retention หรือระยะเวลาที่เก็บ event ไว้ให้อ่านย้อนหลังได้นานขึ้น ต้องการ replay หรือการนำ event ที่เก็บไว้กลับมาประมวลผลอีกครั้ง เช่น คำนวณรายงานด้วยเงื่อนไขใหม่ และต้องการเพิ่ม throughput หรือจำนวน event ที่ระบบรับและประมวลผลได้ต่อวินาที โดยให้แต่ละ service อ่านตามจังหวะของตัวเอง

เมื่อ JetStream เริ่มมีข้อจำกัดสำหรับ event stream ขนาดใหญ่

JetStream ยังสามารถทำงานกับหลาย consumer และเก็บ event แบบ durable ได้ แต่เมื่อ event volume สูงขึ้น เราต้องวางแผนเรื่อง storage, retention, replay และการ scale broker อย่างจริงจัง

ตัวอย่างเช่น Analytics Service อาจหยุดไปหลายชั่วโมงเพื่อ deploy version ใหม่ หลังจากกลับมา online มันต้องอ่าน event ที่เกิดขึ้นระหว่างนั้นทั้งหมด หรือทีม Data อาจต้องสร้าง report ใหม่จาก event ย้อนหลังหลายเดือน

ถ้า event ถูกใช้เป็น source สำหรับหลายระบบ เราต้องการ model ที่แบ่ง stream ออกเป็นส่วน ๆ เพื่อให้ consumer หลายตัวอ่านพร้อมกันได้ เรายังต้องการเครื่องมือเชื่อมข้อมูลกับระบบอื่น เช่น connector สำหรับส่ง event เข้า data warehouse รวมถึง stream processing สำหรับคำนวณจาก event ที่เข้ามาต่อเนื่อง และ data pipeline ที่ส่งข้อมูลผ่านขั้นตอนเหล่านี้

นี่คือ pain point ที่ทำให้เราเริ่มมองหา event streaming platform ที่ออกแบบมาสำหรับ event volume และ consumer จำนวนมาก

Kafka แบ่ง event stream ด้วย topic และ partition

ในระบบของเรา Order Service ส่ง event เข้า topic ชื่อ order.events ให้มอง topic เป็นชื่อที่ producer ใช้ส่งข้อมูล และ consumer ใช้เลือกว่าจะอ่าน event ชุดไหน ภายใน topic นี้ Kafka แบ่งข้อมูลออกเป็น partition เช่น partition 0, 1 และ 2 แต่ละ partition เป็น log ของตัวเอง คือรายการ event ที่เขียนต่อท้ายตามลำดับ

เมื่อมี order ใหม่ event จะถูกเขียนลง partition ใด partition หนึ่ง ไม่ได้คัดลอกลงทั้งสาม partition เพื่อแบ่งงาน การมีหลาย partition ทำให้กระจายงานอ่านและเขียนไปยังหลาย broker หรือ server ที่เก็บข้อมูลของ Kafka ได้ ส่วนการทำสำเนาข้อมูลเพื่อรับมือ broker ล่มเรียกว่า replication ซึ่งเป็นคนละเรื่องกับการแบ่ง partition

ก่อนดูภาพ ให้รู้จัก consumer group อีกคำหนึ่ง: เป็นกลุ่ม consumer ที่ช่วยกันอ่าน topic สำหรับงานเดียวกัน ในตัวอย่างนี้ Payment มี group ของตัวเอง และ Analytics มีอีก group หนึ่ง ทั้งสอง group จึงอ่าน event ชุดเดียวกันได้ โดยภายในแต่ละ group ค่อยแบ่ง partition ให้สมาชิกช่วยกันอ่าน

flowchart LR
    Order["Order Service"] --> Topic["Kafka topic<br/>order.events"]

    subgraph Partitions["Partitions"]
        P0["partition 0"]
        P1["partition 1"]
        P2["partition 2"]
    end

    Topic --> P0
    Topic --> P1
    Topic --> P2

    P0 --> Payment["Payment consumer group"]
    P1 --> Payment
    P2 --> Payment
    P0 --> Analytics["Analytics consumer group"]
    P1 --> Analytics
    P2 --> Analytics

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class Order order
    classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
    class Topic,P0,P1,P2 broker
    classDef payment fill:#fce7f3,stroke:#be185d,color:#0f172a
    class Payment payment
    classDef analytics fill:#ede9fe,stroke:#7c3aed,color:#0f172a
    class Analytics analytics
    style Partitions fill:#fef3c7,stroke:#a16207,color:#0f172a,stroke-width:2px

ถ้า Payment group มี consumer สามตัว Kafka สามารถแบ่งให้แต่ละตัวอ่านหนึ่ง partition ได้ แต่ละ partition จะมี consumer ที่รับผิดชอบเพียงตัวเดียวภายใน group ณ เวลาหนึ่ง ถ้าเพิ่มตัวที่สี่ในตัวอย่างที่มีเพียง topic นี้ มันจะยังไม่มี partition ให้รับผิดชอบ เมื่อสมาชิกเข้าหรือออก Kafka จะจัดสรร partition ใหม่ กระบวนการนี้เรียกว่า rebalance ขณะเดียวกัน Analytics group ก็มีการแบ่งงานและตำแหน่งการอ่านของตัวเอง ดูแนวคิดนี้ใน Kafka consumer groups

Kafka ไม่ได้แทนที่ JetStream ในทุก service

ตอนนี้เราเห็น trade-off ชัดขึ้น Kafka เหมาะเมื่อ event ต้องถูกเก็บและอ่านโดยหลายระบบ ต้องการ replay, partition และ throughput สูง

แต่ถ้า Order API แค่ส่ง FindRider ให้ worker ภายในระบบ การใช้ Kafka อาจเพิ่ม operational cost และ latency โดยไม่จำเป็น JetStream หรือ message broker ที่เล็กกว่ายังเหมาะกว่า

การเลือกจึงควรดู business workflow และ ownership ของแต่ละ service ไม่ใช่เลือก broker เดียวให้ทั้ง system

ในตอนถัดไปเราจะสรุป communication style แต่ละแบบ พร้อม use cases ว่า workflow ไหนควรใช้ synchronous request, asynchronous job, event notification, JetStream หรือ Kafka