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

  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

API ยังไหว แต่ rider job เริ่มค้าง

ช่วงเที่ยงมีหลายร้านรับ order พร้อมกัน API ยังตอบ request ได้ แต่ worker ที่หา rider เริ่มทำ job ไม่ทัน การเพิ่ม application instance ช่วยได้บางส่วน แต่เราต้องเพิ่มทั้ง API และ worker ไปด้วยกัน ทั้งที่ส่วนที่ต้องการ scale จริง ๆ คือ rider job

จากตอนที่ 2 เรารู้แล้วว่าแต่ละ instance มี local queue ของตัวเอง หาก instance หนึ่งมี job ค้าง worker ของอีก instance ก็เข้ามาช่วยอ่าน queue นั้นไม่ได้

เราจึงแยก Rider Matching worker ไปเป็นอีก process เพื่อเพิ่มหรือลดจำนวน worker ได้อิสระจาก API คราวนี้ request handler จะฝาก order ID ไว้ใน channel เดิมไม่ได้แล้ว เพราะ worker อยู่คนละ memory space

ก่อนเลือกเครื่องมือ: เราต้องการ solution แบบไหน?

ปัญหาของเราคือ Order API และ Rider Matching worker อยู่คนละ process ดังนั้น local queue ที่อยู่ใน memory ของ Order API จึงส่ง job ให้ worker โดยตรงไม่ได้

สิ่งที่เราต้องการคือ shared message channel ที่อยู่ระหว่าง process Order API ส่ง message เข้า channel นี้โดยไม่ต้องรู้ว่า worker มีกี่ตัว หรือแต่ละตัวอยู่ที่ไหน Rider Matching worker ทุกตัวเชื่อมต่อเข้ามาที่ channel เดียวกัน เมื่อมี message เข้ามา ระบบจะเลือก worker หนึ่งตัวให้รับ job ไปทำ

สีใน diagram ใช้เหมือนกันตลอดชุดบทความ: Order / Order API = สีฟ้า; Rider Matching = สีเขียว; Queue / broker / event stream = สีเหลือง สีบอก service หรือประเภทขององค์ประกอบ โดยชื่อในกล่องระบุหน้าที่ ส่วนกรอบแดงเส้นประพร้อมคำว่า “stopped” บอก worker ที่หยุดทำงาน

flowchart LR
    API["Order API"] --> Channel["Shared message channel"]
    Channel --> W1["Rider Matching worker 1"]
    Channel --> W2["Rider Matching worker 2"]
    Channel --> W3["Rider Matching worker 3"]

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class API order
    classDef broker fill:#fef3c7,stroke:#a16207,color:#0f172a
    class Channel broker
    classDef rider fill:#dcfce7,stroke:#15803d,color:#0f172a
    class W1,W2,W3 rider

แนวคิดนี้แยกความรับผิดชอบออกจากกัน Order API รับผิดชอบการส่ง job Worker รับผิดชอบการประมวลผล job ส่วน shared channel รับผิดชอบการ route message ไปยัง worker ที่เชื่อมต่ออยู่

Order API จึงไม่ต้องมีรายชื่อ worker ใน code เราเพิ่มหรือลด worker ได้โดยไม่ต้องแก้ publisher และเมื่อ worker ตัวหนึ่งหยุดทำงาน worker ตัวอื่นยังสามารถรับ message ใหม่ได้ หาก shared channel รองรับ

ตัวกลางที่ทำหน้าที่รับและส่งต่อ message ผ่าน network เรียกว่า message broker เราลองใช้ Core NATS มาสร้าง shared message channel นี้กัน

ลองสร้าง shared message channel ด้วย Core NATS

ใน NATS, subject คือชื่อของ message channel เราตั้งชื่อ channel สำหรับ job นี้เป็น rider.find

Order API คือ publisher มัน publish FindRider พร้อม order ID ไปที่ subject rider.find

Rider Matching worker คือ subscriber มัน subscribe subject rider.find เพื่อรอ message ใหม่

flowchart LR
    API["Publisher: Order API"] -->|publish FindRider| NATS["Core NATS subject: rider.find"]
    NATS -->|deliver message| W1["Subscriber: Matching worker 1"]
    NATS -->|deliver message| W2["Subscriber: Matching worker 2"]

    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

Message ของเราคือ FindRider พร้อม order ID เป็น command เพราะขอให้ subscriber ทำ job เฉพาะอย่าง การมี broker ไม่ได้เปลี่ยน command ให้กลายเป็น event Publisher ยังต้องการให้มี worker รับผิดชอบหา rider ให้ order นั้น

Core NATS ทำหน้าที่เป็น communication layer ระหว่าง process มันช่วย route message จาก publisher ไปยัง subscriber ที่กำลังเชื่อมต่ออยู่ รายละเอียดเรื่อง message durability และกรณี worker offline จะเห็นชัดขึ้นเมื่อเราลองใช้งานจริงท้ายบท (อ้างอิง: NATS — Core NATS)

หลาย worker ต้องแบ่ง job ไม่ใช่หา rider ซ้ำทุกตัว

หาก worker ทุกตัว subscribe แบบปกติ แต่ละตัวจะได้รับสำเนา message เราจึงใช้ queue group ชื่อ rider-matchers ให้ NATS เลือกสมาชิกหนึ่งตัวในกลุ่มที่ subscribe อยู่เพื่อรับแต่ละ message

flowchart LR
    API1["Order API 1"] --> NATS["Core NATS: rider.find"]
    API2["Order API 2"] --> NATS
    subgraph Group["Queue group: rider-matchers"]
        W1["Matching worker 1"]
        W2["Matching worker 2"]
    end
    NATS -->|เลือกหนึ่งสมาชิกต่อ message| W1
    NATS -->|เลือกหนึ่งสมาชิกต่อ message| W2

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class API1,API2 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 Group fill:#dcfce7,stroke:#15803d,color:#0f172a,stroke-width:2px

Worker ที่เพิ่มเข้ามาต้องใช้ทั้ง subject และชื่อ group เดียวกัน หากใช้คนละ group แต่ละ group จะได้รับสำเนาแยกกัน การกระจายนี้ไม่ได้รับประกันว่าเลือก worker ที่ว่างที่สุดหรือวัดเวลาทำ job จริงให้เรา ชื่อ queue group ก็ไม่ได้แปลว่ามี durable queue เก็บ job อยู่เบื้องหลัง

ตัวอย่าง Go: เปลี่ยนจาก local queue เป็น subject

สมมติว่ามี NATS server รันที่ localhost:4222 แล้ว สร้างโฟลเดอร์ทดลองด้วย go mod init example.com/rider-messaging และเพิ่ม client ด้วย go get github.com/nats-io/nats.go โค้ดนี้ใช้ไฟล์ main.go เดียว เลือกว่าจะทำหน้าที่เป็น worker หรือ publisher จาก argument:

package main

import (
	"fmt"
	"log"
	"os"
	"time"

	"github.com/nats-io/nats.go"
)

func run() error {
	if len(os.Args) < 2 {
		return fmt.Errorf("usage: worker | send ORDER_ID")
	}
	nc, err := nats.Connect(nats.DefaultURL)
	if err != nil {
		return err
	}
	defer nc.Close()

	switch os.Args[1] {
	case "worker":
		sub, err := nc.QueueSubscribeSync("rider.find", "rider-matchers")
		if err != nil {
			return err
		}
		if err := nc.FlushTimeout(2 * time.Second); err != nil {
			return err
		}
		fmt.Println("worker พร้อมรับ job")
		for {
			msg, err := sub.NextMsg(time.Second)
			if err == nats.ErrTimeout {
				continue
			}
			if err != nil {
				return err
			}
			fmt.Printf("worker %d จำลองหา rider ให้ %s ", os.Getpid(), msg.Data)
		}
	case "send":
		if len(os.Args) != 3 {
			return fmt.Errorf("usage: send ORDER_ID")
		}
		if err := nc.Publish("rider.find", []byte(os.Args[2])); err != nil {
			return err
		}
		return nc.FlushTimeout(2 * time.Second)
	default:
		return fmt.Errorf("unknown mode: %s", os.Args[1])
	}
}

func main() {
	if err := run(); err != nil {
		log.Fatal(err)
	}
}

เปิดสอง terminal รัน go run . worker ทั้งคู่ และรอให้แสดงว่าพร้อมรับ job จากนั้นใช้ terminal ที่สามรัน go run . send order-1 และ go run . send order-2 แต่ละ message จะไปสมาชิกหนึ่งตัวใน group โดยไม่จำเป็นต้องสลับกันพอดี

จุดที่เปลี่ยนจากตอนก่อนคือ publisher ส่งไปที่ rider.find แทน channel ส่วน worker รับผ่าน subscription การหา rider จริงยังเป็น business logic เดิมที่เราเพียงย้ายที่รัน ตัวอย่าง print log แทนการทำ job นั้น เพื่อให้เห็นเส้นทาง message ชัด ๆ

FlushTimeout รอให้ server ตอบกลับหลังรับคำสั่งก่อนหน้า ไม่ได้ยืนยันว่า worker ทำ job สำเร็จหรือ message ถูก save ลง disk ส่วน QueueSubscribeSync หมายถึง client อ่าน message เอง ไม่ได้ทำให้ publisher รอผลการหา rider (อ้างอิง API: nats.go)

เราแยกการ scale ได้ แต่เพิ่มสิ่งที่ต้องดูแล

flowchart LR
    API1["Order API 1"] --> NATS["Broker: Core NATS<br/>subject: rider.find"]
    API2["Order API 2"] --> NATS
    subgraph Group["Queue group: rider-matchers<br/>scale workers independently"]
        W1["Matching worker 1"]
        W2["Matching worker 2"]
    end
    NATS -->|deliver one message to one worker| W1
    NATS -->|deliver one message to one worker| W2

    classDef order fill:#dbeafe,stroke:#2563eb,color:#0f172a
    class API1,API2 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 Group fill:#dcfce7,stroke:#15803d,color:#0f172a,stroke-width:2px

ตอนนี้ scale worker ได้โดยไม่ต้อง scale API และ API ไม่ต้อง track รายชื่อ worker เอง นี่คือประโยชน์หลักของการแยก process ออกจากกัน

แต่การแยก process ทำให้เราต้องดูแลสิ่งที่เคยซ่อนอยู่ใน function call ด้วย เราต้องดูแล broker, connection และกรณี network ขาดหาย เราต้อง track ว่า worker ยังเชื่อมต่ออยู่หรือไม่ และอ่าน message ได้ทันหรือเปล่า

หาก worker รับ message ได้เร็ว แต่ใช้เวลานานในการหา rider หรือบันทึกลง database message ถัดไปอาจถูกส่งเข้ามารอใน client หรือ connection buffer ของ worker ตัวนั้น

Core NATS ไม่ได้รู้ว่า worker ทำ job เสร็จแล้วหรือยัง จึงไม่ได้หยุดส่ง message เพียงเพราะ worker ยัง process message ก่อนหน้าไม่เสร็จ และ queue group ก็ไม่ได้วัดว่า worker ตัวไหนทำงานช้ากว่ากันเพื่อจัดสรร message ใหม่ให้เหมาะสม

ผลคือ worker ที่ช้าอาจมี message ค้างใน buffer ขณะที่ worker ตัวอื่นยังมี capacity เหลือ การเพิ่ม worker ช่วยเพิ่ม throughput ได้เมื่อมี capacity เหลือ แต่ไม่ได้ทำให้ระบบเร็วขึ้นไม่จำกัด เพราะ worker ยังต้องใช้ database, map service และ service อื่น ๆ ร่วมกัน

แล้วถ้าปิด worker ทุกตัวก่อนส่ง job ล่ะ?

ลองหยุด worker ทั้งสองตัว แล้วส่ง go run . send order-3 จากนั้นเปิด worker กลับมาใหม่ จะพบว่าไม่มี job order-3 รออยู่

ตอนแรก Core NATS ดูเหมือนจะแก้ปัญหาของเราได้ครบ Order API ไม่ต้องรู้จัก worker แต่ละตัว เราสามารถเพิ่ม worker เพื่อช่วยกันรับ job ได้ และ publisher กับ subscriber ก็แยก deploy กันได้

แต่เมื่อไม่มี worker รับ message ในเวลาที่ส่ง เราก็เห็นข้อจำกัดที่สำคัญ Core NATS ไม่เก็บ message รอ subscriber ที่ offline ถ้าไม่มี subscriber เชื่อมต่ออยู่ตอน publisher ส่ง message ระบบจะไม่เก็บไว้ใน memory หรือ disk เพื่อส่งภายหลัง

นอกจากนี้ Core NATS ไม่มี consumer acknowledgement สำหรับบอกว่า job ทำสำเร็จแล้ว ถ้า worker รับ message ไปแล้ว process crash ระหว่างหา rider ระบบก็ไม่รู้ว่าต้องส่ง message เดิมให้ worker ตัวอื่นอีกครั้ง

ลักษณะ delivery นี้เรียกว่า at-most-once delivery message จะถูกส่งให้ subscriber ที่พร้อมรับในเวลานั้นมากที่สุดหนึ่งครั้ง (อ้างอิง: NATS — Core NATS is ephemeral)

รูปแบบนี้เหมาะกับข้อมูลที่ยอมให้สูญหายได้ หรือข้อมูลใหม่สามารถแทนที่ข้อมูลเก่าได้ เช่น live location หรือ cache invalidation แต่ FindRider เป็น job ที่หายไม่ได้ ถ้า message หาย order อาจไม่มี rider มารับ

เราจึงแก้ปัญหา ส่ง job ข้าม process ได้แล้ว แต่ยังไม่ได้แก้ปัญหา เก็บ job ไว้จน subscriber พร้อม และรู้ว่า job ไหนต้องนำกลับมาทำต่อ ตอนถัดไปเราจะเพิ่มความสามารถด้าน durability เพื่อให้ message อยู่รอดหลัง worker offline หรือ process restart