บทที่ 11 · Part 3 — Production & Advanced
Production & Observability
Worker tuning, metrics, logging, Search Attributes, Priority & Fairness สำหรับ multi-tenant และ payload limits
Workflow ที่ทำงานถูกต้องบนเครื่อง dev กับระบบที่รับธุรกรรมจริงวันละ 1 ล้านรายการ เป็นคนละเรื่องกัน บทนี้ว่าด้วยสิ่งที่ต้องเตรียมก่อนเปิดให้ลูกค้าใช้จริง — ตั้ง worker ให้พอดี, มองเห็นสิ่งที่เกิดขึ้น, และไม่ให้ร้านค้ารายใหญ่กินคิวจนรายเล็กไม่ได้รัน
จบบทนี้คุณจะ
- ตั้งค่า worker ตามภาระจริง และรู้ว่าต้องดู metric ตัวไหน
- อ่าน metric ที่สำคัญและตั้ง alert ที่มีความหมาย
- เขียน log ที่ replay-safe และไม่ละเมิด PDPA/PCI
- ใช้ Search Attributes ให้ ops ค้นธุรกรรมเจอ
- ใช้ Priority & Fairness แก้ปัญหา multi-tenant starvation
- หลบข้อจำกัดเรื่องขนาด payload และ history
Worker tuning
Worker คือ process ของเรา ที่ poll งานจาก task queue มาทำ ค่า default เหมาะกับการเริ่มต้นแต่ไม่เหมาะกับ production
w := worker.New(c, "wallet-core", worker.Options{
// จำนวน activity ที่รันพร้อมกันได้สูงสุดใน process นี้ (default 1000)
MaxConcurrentActivityExecutionSize: 200,
// จำนวน workflow task ที่ประมวลผลพร้อมกัน (default 1000)
MaxConcurrentWorkflowTaskExecutionSize: 200,
// จำนวน poller (default 2 ทั้งคู่)
MaxConcurrentActivityTaskPollers: 8,
MaxConcurrentWorkflowTaskPollers: 8,
// จำกัด activity ที่ยิงออกต่อวินาทีทั้ง worker — กันถล่ม API ธนาคาร
WorkerActivitiesPerSecond: 50,
// รอให้งานที่ทำอยู่จบก่อนปิด — สำคัญมากตอน rolling deploy
WorkerStopTimeout: 30 * time.Second,
})
อย่าตั้งค่าเดาสุ่ม ให้ metric เป็นคนบอก
ตัวชี้วัดหลักคือ temporal_activity_schedule_to_start_latency —
เวลาที่ task นั่งรอในคิวก่อนมี worker มาหยิบ
- ค่าสูงขึ้นเรื่อยๆ = worker ไม่พอ → เพิ่ม pollers หรือเพิ่มจำนวน worker
- ค่าต่ำแต่ CPU เต็ม = worker ทำงานหนักเกิน → ลด concurrency ลงหรือ scale out
- ค่าต่ำและ CPU ว่าง = ตั้งค่าเหมาะสมแล้ว
แยก task queue ตามลักษณะงาน
อย่าเอา activity ที่คุยกับธนาคาร (ช้า, มี rate limit) ไปอยู่คิวเดียวกับงานที่เร็ว เพราะงานช้าจะกิน slot จนงานเร็วอด
// worker สำหรับ orchestration — เบา เร็ว
walletWorker := worker.New(c, "wallet-core", worker.Options{
MaxConcurrentWorkflowTaskExecutionSize: 500,
})
walletWorker.RegisterWorkflow(workflows.TopUpWorkflow)
// worker สำหรับงานคุยธนาคาร — จำกัด rate ตามโควตาที่ธนาคารให้
bankWorker := worker.New(c, "bank-io", worker.Options{
MaxConcurrentActivityExecutionSize: 20,
WorkerActivitiesPerSecond: 30,
})
bankWorker.RegisterActivity(&activities.BankActivities{})
แล้วชี้ activity ไปคิวที่ถูกต้องจากใน workflow:
bankCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
TaskQueue: "bank-io",
StartToCloseTimeout: 60 * time.Second,
})
Sticky cache คือเหตุผลที่ worker ไม่ควร restart บ่อย
Worker เก็บ workflow ที่กำลังทำงานไว้ใน memory (sticky execution)
ทำให้ไม่ต้อง replay ใหม่ทุกครั้ง ถ้า temporal_sticky_cache_miss สูง
แปลว่า workflow ถูกเตะออกจาก cache บ่อย — worker restart ถี่เกินไป
หรือ WorkflowCacheSize เล็กเกินไป ผลคือ CPU ถูกใช้ไปกับการ replay ซ้ำๆ
Metrics ที่ต้องดู
ต่อ metrics handler เข้ากับ client ตอนสร้าง:
c, err := client.Dial(client.Options{
MetricsHandler: sdktally.NewMetricsHandler(newPrometheusScope(prometheus.Configuration{
ListenAddress: "0.0.0.0:9090",
TimerType: "histogram",
})),
})
| Metric | บอกอะไร | ตั้ง alert เมื่อ |
|---|---|---|
temporal_workflow_task_execution_failed | workflow task ล้ม | มี tag failure_reason="NonDeterminismError" แม้แค่ครั้งเดียว |
temporal_activity_schedule_to_start_latency | งานรอคิวนานแค่ไหน | p99 สูงขึ้นต่อเนื่อง = worker ไม่พอ |
temporal_workflow_task_schedule_to_start_latency | workflow task รอคิว | สูง = workflow worker ไม่พอ |
temporal_activity_execution_failed | activity ล้ม | อัตราสูงผิดปกติ = ระบบปลายทางมีปัญหา |
temporal_workflow_endtoend_latency | ธุรกรรมใช้เวลารวมเท่าไร | เกิน SLA ที่ให้ลูกค้าไว้ |
temporal_worker_task_slots_available | slot ว่างเหลือเท่าไร | ใกล้ 0 ตลอด = ต้อง scale |
temporal_sticky_cache_miss | replay เกิดบ่อยแค่ไหน | พุ่งขึ้น = worker ไม่เสถียร |
alert ตัวที่สำคัญที่สุดคือ NonDeterminismError
temporal_workflow_task_execution_failed มี tag failure_reason ที่แยกได้ว่า
เป็น NonDeterminismError, GrpcMessageTooLarge หรือ WorkflowError
ตัวแรกแปลว่ามี deploy ที่ทำ workflow ที่ค้างอยู่พัง (บทที่ 10)
ต้องรู้ภายในไม่กี่นาที ไม่ใช่รอลูกค้าโทรมาถาม
[!NOTE] หน่วยของ histogram ไม่เหมือนกันทุก SDK Go และ Java ใช้วินาที ส่วน SDK ที่สร้างบน Core (TypeScript, Python, .NET, Ruby) ใช้มิลลิวินาที ถ้าทีมมีหลายภาษาในระบบเดียวกัน dashboard ต้องแปลงหน่วยให้ตรง ไม่งั้นกราฟจะเพี้ยนไป 1,000 เท่า
Logging
ใช้ logger ของ SDK เท่านั้นใน workflow — มันรู้ว่ากำลัง replay อยู่จึงไม่ log ซ้ำ
func TopUpWorkflow(ctx workflow.Context, req TopUpRequest) (TopUpResult, error) {
// ติด field ที่ต้องใช้ทุกบรรทัดไว้ครั้งเดียว
logger := log.With(workflow.GetLogger(ctx),
"txnID", req.TransactionID,
"walletID", req.WalletID,
)
logger.Info("top-up started", "amountSatang", req.AmountSatang)
// ...
}
ต่อกับ slog เพื่อให้ log ของ SDK ไปรวมกับ log ของแอปเรา:
handler := slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo})
c, err := client.Dial(client.Options{
Logger: log.NewStructuredLogger(slog.New(handler)),
})
ห้าม log ข้อมูลเหล่านี้ลง log ธรรมดาเด็ดขาด
- เลขบัตร (PAN), CVV, expiry — ข้อกำหนด PCI-DSS
- OTP, รหัสผ่าน, token, API key
- เลขบัตรประชาชน, เลขบัญชีเต็ม — PDPA
ใช้ค่าที่อ้างอิงได้แทน: transaction ID, wallet ID, bank reference ถ้าจำเป็นต้องเห็นเลขบัญชีให้ mask เหลือ 4 ตัวท้าย และจำไว้ว่า input/output ของ workflow ถูกเก็บใน event history ทั้งหมด ซึ่งเปิดดูได้จาก Web UI — การเข้ารหัสส่วนนี้เป็นเรื่องของบทที่ 12
fmt.Println และ log.Println ใน workflow จะพิมพ์ซ้ำทุกครั้งที่ replay
ทำให้ log ปลอมและอ่านไม่รู้เรื่อง — ถ้าจำเป็นจริงๆ ให้เช็ค workflow.IsReplaying(ctx) ก่อน
Search Attributes
Web UI ค้นด้วย Workflow ID ได้อยู่แล้ว แต่ ops มักถามคำถามแบบ "ขอรายการของร้าน M-042 ที่ค้างเกิน 1 ชั่วโมง" ซึ่งต้องใช้ Search Attributes
สร้าง attribute ใน namespace ก่อน (ทำครั้งเดียว):
temporal operator search-attribute create --name MerchantId --type Keyword
temporal operator search-attribute create --name TxnStage --type Keyword
temporal operator search-attribute create --name AmountSatang --type Int
ประกาศ key แบบมี type แล้วใช้ได้ทั้งฝั่ง client และ workflow:
var (
MerchantIDKey = temporal.NewSearchAttributeKeyKeyword("MerchantId")
TxnStageKey = temporal.NewSearchAttributeKeyKeyword("TxnStage")
AmountSatangKey = temporal.NewSearchAttributeKeyInt64("AmountSatang")
)
// ตอนเริ่ม workflow
_, err := c.ExecuteWorkflow(ctx, client.StartWorkflowOptions{
ID: "topup-" + txnID,
TaskQueue: "wallet-core",
TypedSearchAttributes: temporal.NewSearchAttributes(
MerchantIDKey.ValueSet("M-042"),
AmountSatangKey.ValueSet(100_00),
),
}, workflows.TopUpWorkflow, req)
// อัปเดตระหว่างทางจากใน workflow
if err := workflow.UpsertTypedSearchAttributes(ctx,
TxnStageKey.ValueSet("awaiting_otp")); err != nil {
return TopUpResult{}, err
}
แล้ว ops ก็ค้นได้:
temporal workflow list --query \
'MerchantId = "M-042" AND TxnStage = "awaiting_otp" AND ExecutionStatus = "Running"'
Search Attribute กับ Memo ต่างกันตรงค้นได้หรือไม่ได้
- Search Attribute — index ไว้ ค้นได้ แต่ต้องประกาศ type ก่อน และมีจำนวนจำกัด
- Memo — แนบข้อมูลไปเฉยๆ ค้นไม่ได้ แต่ใส่อะไรก็ได้ ไม่ต้องประกาศ
ใช้ Search Attribute เฉพาะ field ที่ ops ต้องใช้ค้นจริงๆ อย่ายัดทุกอย่างลงไปเพราะมันกินพื้นที่ visibility store และ ห้ามใส่ข้อมูลส่วนบุคคล เพราะ search attribute ไม่ถูกเข้ารหัสด้วย codec
Priority & Fairness — ปัญหา multi-tenant
ระบบ wallet ที่ให้บริการหลายร้านค้าจะเจอปัญหานี้แน่นอน: ร้านใหญ่ยิงงานเข้ามา 100,000 รายการตอนปิดรอบ ส่วนร้านเล็กส่งมา 10 รายการ — ด้วยคิวแบบ FIFO ธรรมดา ร้านเล็กต้องรอจนกว่างานของร้านใหญ่จะหมด
Fairness แก้ตรงนี้ด้วยการให้แต่ละ key มี virtual queue ของตัวเอง แล้ว dispatch สลับกันไป ร้านเล็กจึงได้รันแทรกเสมอ
_, err := c.ExecuteWorkflow(ctx, client.StartWorkflowOptions{
ID: "topup-" + txnID,
TaskQueue: "wallet-core",
Priority: temporal.Priority{
PriorityKey: 1, // 1–5, เลขน้อย = สำคัญกว่า, default 3
FairnessKey: merchantID, // 1 key = 1 virtual queue
FairnessWeight: 2.0, // น้ำหนักเทียบกับ key อื่น, default 1.0
},
}, workflows.TopUpWorkflow, req)
Priority เป็นคนละเรื่อง — มันจัดลำดับความสำคัญของประเภทงาน ในคิวเดียวกัน
| Priority | Fairness | |
|---|---|---|
| แก้ปัญหา | งานด่วนต้องมาก่อนงาน batch | ผู้ใช้รายใหญ่ไม่ให้กินคิวรายเล็ก |
| ค่า | integer 1–5 (default 3) | string key + weight (default 1.0) |
| กลไก | sub-queue ตามลำดับความสำคัญ | round-robin ข้าม virtual queue |
| ใช้ร่วมกัน | Priority เลือก sub-queue ก่อน แล้ว Fairness ทำงานภายในแต่ละระดับ |
ตัวอย่างการแบ่งชั้นบริการ:
| Fairness Key | Weight | สัดส่วนที่ได้รัน |
|---|---|---|
tier-premium | 5.0 | 50% |
tier-standard | 3.0 | 30% |
tier-free | 2.0 | 20% |
จำกัด rate ต่อคิวและต่อ key ได้ด้วย:
temporal task-queue config set \
--task-queue bank-io \
--task-queue-type activity \
--queue-rps-limit 500 \
--queue-rps-limit-reason "โควตารวมที่ธนาคารให้" \
--fairness-key-rps-limit-default 20 \
--fairness-key-rps-limit-reason "กันร้านเดียวยิงจนเต็มโควตา"
ข้อควรรู้ก่อนใช้ Fairness
- ทั้ง Priority และ Fairness อยู่ในสถานะ Public Preview — Priority ใช้ฟรี ส่วน Fairness เป็นฟีเจอร์เสียเงินบน Temporal Cloud ต้องเช็คสัญญาก่อนออกแบบพึ่งพามัน
- บน self-hosted ต้องเปิด dynamic config:
matching.useNewMatcher,matching.enableFairnessและmatching.enableMigration - Fairness ทำงานตอน dispatch ไม่ได้ดึงงานที่ worker กำลังทำอยู่กลับมา ช่วงแรกหลังเปิดใช้จึงอาจยังเห็นร้านใหญ่ครองอยู่
- น้ำหนักมีผลตอน schedule ไม่ใช่ตอน dispatch — เปลี่ยน weight ไม่ได้จัดลำดับงานที่ค้างอยู่ในคิวใหม่
- ถ้าไม่มี backlog เลย (งานถูกหยิบไปทำทันทีตลอด) ก็ไม่ต้องใช้ Fairness
ข้อจำกัดเรื่องขนาด
| ขีดจำกัด (ค่า default) | ผลเมื่อเกิน |
|---|---|
| 2MB ต่อ 1 payload | activity หรือ workflow ล้ม |
| 4MB ต่อ 1 gRPC message | ส่งไม่ได้ — เห็นเป็น GrpcMessageTooLarge |
| 50MB ต่อ history (ควรอยู่ต่ำกว่า 10MB) | workflow ถูกบังคับให้จบ |
ค่าเหล่านี้ปรับได้ระดับ cluster แต่อย่าไปแก้เพื่อหลบปัญหา เพราะรากของปัญหา คือการส่งข้อมูลก้อนใหญ่เข้าออก workflow ซึ่งทำให้ replay ช้าลงทุกครั้งด้วย
วิธีที่ถูกคือส่ง reference แทนข้อมูล:
// อย่าทำ — statement ทั้งเดือนเป็น MB เข้าไปอยู่ใน history ตลอดกาล
var statement []Transaction
_ = workflow.ExecuteActivity(actCtx, a.FetchStatement, bank, day).Get(actCtx, &statement)
// ทำแบบนี้ — activity เก็บลง object storage แล้วคืนแค่ key
var ref StatementRef // { Bucket, Key, RowCount, Checksum }
_ = workflow.ExecuteActivity(actCtx, a.FetchStatementToStorage, bank, day).Get(actCtx, &ref)
_ = workflow.ExecuteActivity(actCtx, a.ReconcileFromStorage, ref).Get(actCtx, nil)
ส่วน history ที่โตจากจำนวนรอบการทำงาน แก้ด้วย Continue-As-New ตามบทที่ 8
ตารางแก้ปัญหาหน้างาน
| อาการ | สาเหตุที่พบบ่อย | ทำอะไร |
|---|---|---|
| Workflow ค้าง Running แต่ไม่เดิน | non-determinism หลัง deploy | ดู failure_reason ใน metric แล้ว rollback |
| ธุรกรรมช้าผิดปกติทั้งระบบ | worker ไม่พอ | ดู schedule_to_start_latency แล้วเพิ่ม worker |
| ร้านเล็กบ่นว่ารายการค้าง | คิวถูกร้านใหญ่ครอง | เปิด Fairness ด้วย merchant ID |
| Activity retry ไม่หยุด | error ที่ควรเป็น non-retryable | แก้การจัดประเภท error (บทที่ 6) |
| CPU worker สูงแต่งานไม่เดิน | replay ซ้ำจาก cache miss | ลดการ restart worker, เพิ่ม cache size |
GrpcMessageTooLarge | payload ใหญ่เกิน | เปลี่ยนไปส่ง reference |
Checklist ก่อนเปิดให้ลูกค้าใช้
- Worker แยก task queue ตามลักษณะงานแล้ว (orchestration / งานคุยระบบภายนอก)
- มี metrics endpoint และ dashboard ที่มี schedule-to-start latency เป็นตัวหลัก
- มี alert บน
failure_reason="NonDeterminismError"ที่ดังทันที - Log ผ่าน
workflow.GetLoggerทั้งหมด และไม่มีข้อมูลอ่อนไหวหลุด - Search Attributes ครอบคลุมคำถามที่ ops ถามบ่อย
- Multi-tenant ตั้ง fairness key ตาม merchant แล้ว
- ไม่มี payload ก้อนใหญ่ไหลผ่าน workflow
WorkerStopTimeoutตั้งไว้ให้ rolling deploy ไม่ตัดงานกลางคัน