บทที่ 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_failedworkflow task ล้มมี tag failure_reason="NonDeterminismError" แม้แค่ครั้งเดียว
temporal_activity_schedule_to_start_latencyงานรอคิวนานแค่ไหนp99 สูงขึ้นต่อเนื่อง = worker ไม่พอ
temporal_workflow_task_schedule_to_start_latencyworkflow task รอคิวสูง = workflow worker ไม่พอ
temporal_activity_execution_failedactivity ล้มอัตราสูงผิดปกติ = ระบบปลายทางมีปัญหา
temporal_workflow_endtoend_latencyธุรกรรมใช้เวลารวมเท่าไรเกิน SLA ที่ให้ลูกค้าไว้
temporal_worker_task_slots_availableslot ว่างเหลือเท่าไรใกล้ 0 ตลอด = ต้อง scale
temporal_sticky_cache_missreplay เกิดบ่อยแค่ไหนพุ่งขึ้น = 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 เป็นคนละเรื่อง — มันจัดลำดับความสำคัญของประเภทงาน ในคิวเดียวกัน

PriorityFairness
แก้ปัญหางานด่วนต้องมาก่อนงาน batchผู้ใช้รายใหญ่ไม่ให้กินคิวรายเล็ก
ค่าinteger 1–5 (default 3)string key + weight (default 1.0)
กลไกsub-queue ตามลำดับความสำคัญround-robin ข้าม virtual queue
ใช้ร่วมกันPriority เลือก sub-queue ก่อน แล้ว Fairness ทำงานภายในแต่ละระดับ

ตัวอย่างการแบ่งชั้นบริการ:

Fairness KeyWeightสัดส่วนที่ได้รัน
tier-premium5.050%
tier-standard3.030%
tier-free2.020%

จำกัด 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 payloadactivity หรือ 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
GrpcMessageTooLargepayload ใหญ่เกินเปลี่ยนไปส่ง 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 ไม่ตัดงานกลางคัน