บทที่ 13 · Part 3 — Production & Advanced
Advanced Patterns
Local Activity, Polling, Async Activity Completion, Session, Nexus และ Standalone Activities
เครื่องมือในบทก่อนๆ ครอบคลุมงานส่วนใหญ่ได้แล้ว บทนี้เก็บของที่เหลือ — pattern ที่ใช้ไม่บ่อยแต่พอถึงเวลาต้องใช้แล้วไม่มีอะไรทดแทนได้
จบบทนี้คุณจะ
- รู้ว่าเมื่อไหร่ Local Activity คุ้มความเสี่ยง
- เลือก pattern polling ให้ตรงกับความถี่ที่ต้องการ
- รับผลลัพธ์จากระบบภายนอกด้วย Async Activity Completion
- รู้จัก Session, Nexus และ Standalone Activities ว่าแก้ปัญหาอะไร
Local Activity
Activity ปกติต้องผ่าน task queue: workflow ออก command → server บันทึก → worker poll → ทำงาน → รายงานผล → server บันทึก รวมแล้วมี overhead หลายสิบมิลลิวินาที Local Activity ข้ามขั้นตอนนั้น รันในกระบวนการเดียวกับ workflow task เลย
lao := workflow.LocalActivityOptions{
StartToCloseTimeout: 5 * time.Second,
RetryPolicy: &temporal.RetryPolicy{
MaximumAttempts: 3,
},
}
ctx = workflow.WithLocalActivityOptions(ctx, lao)
var tier string
err := workflow.ExecuteLocalActivity(ctx, a.LookupMerchantTier, merchantID).Get(ctx, &tier)
| Activity ปกติ | Local Activity | |
|---|---|---|
| Latency | สูงกว่า (ผ่าน task queue) | ต่ำ |
| เห็นใน history ระหว่างรัน | เห็น (ActivityTaskScheduled) | ไม่เห็นจนกว่าจะจบ |
| Retry ข้าม worker ได้ | ได้ | ไม่ได้ — worker ตาย งานเริ่มใหม่หมด |
| Heartbeat | ได้ | ไม่ได้ |
| Timeout ยาวๆ | ได้ | ไม่เหมาะ |
| Cancel ระหว่างทาง | ได้ | ไม่ได้ |
Local Activity ไม่ทนทานเท่า activity ปกติ
มันไม่มี record ใน history จนกว่าจะทำเสร็จ ถ้า worker ตายกลางคัน งานจะถูกเริ่มใหม่ทั้งหมดตอน replay ห้ามใช้กับอะไรก็ตามที่ทำให้เงินเคลื่อนที่ ต่อให้เร็วกว่าแค่ไหนก็ตาม
ใช้ได้กับ: อ่าน config, lookup ตารางอ้างอิง, คำนวณที่ต้องใช้ข้อมูลภายนอกเล็กน้อย, ตรวจ feature flag — งานที่ทำซ้ำแล้วไม่มีผลข้างเคียง และเสร็จในไม่กี่วินาที
[!TIP] เกณฑ์ตัดสินสั้นๆ ถามว่า "ถ้างานนี้ถูกรันซ้ำ 2 รอบโดยไม่มีใครรู้ จะเสียหายไหม" ถ้าไม่เสียหายและมันเร็วมาก → Local Activity ได้ ถ้าเสียหาย → Activity ปกติเท่านั้น
Polling
โจทย์ที่พบบ่อย: ธนาคารไม่มี webhook ต้อง poll สถานะเอง
วิธีที่ผิดคือวนลูปใน workflow แล้ว Sleep ทุกรอบ เพราะทุกรอบเพิ่ม event 2 ตัวใน history
// อย่าทำ — poll 5 วินาทีเป็นเวลา 6 ชั่วโมง = 8,640 event
for {
var status string
_ = workflow.ExecuteActivity(actCtx, a.PollBankStatus, ref).Get(actCtx, &status)
if status != "pending" {
break
}
_ = workflow.Sleep(ctx, 5*time.Second)
}
Poll ถี่ — วนลูปในตัว activity
ย้ายลูปเข้าไปอยู่ใน activity แล้ว heartbeat ทุกรอบ วิธีนี้ history ได้แค่ event ของ activity 1 ตัว
func (a *Activities) PollUntilSettled(ctx context.Context, ref BankTransferRef) (string, error) {
for {
status, err := a.Bank.GetStatus(ctx, ref)
if err == nil && status != "pending" {
return status, nil
}
// heartbeat ทำ 2 หน้าที่: บอกว่ายังไม่ตาย และรับรู้การ cancel
activity.RecordHeartbeat(ctx, "polling")
if ctx.Err() != nil {
return "", ctx.Err()
}
select {
case <-time.After(5 * time.Second):
case <-ctx.Done():
return "", ctx.Err()
}
}
}
actCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: 6 * time.Hour,
// ต้องสั้นกว่า StartToCloseTimeout เสมอ ไม่งั้น heartbeat ไม่มีความหมาย
HeartbeatTimeout: 15 * time.Second,
})
`time.After` กับ `select` ใช้ได้ใน activity
ข้อห้ามเรื่อง determinism บังคับใช้กับ workflow เท่านั้น
ใน activity เป็นโค้ด Go ปกติทุกอย่าง — ใช้ goroutine, channel, timer ได้ตามใจ
ขอแค่เคารพ ctx.Done()
Poll ห่าง — ใช้ retry เป็นตัวจับเวลา
ถ้า poll ทุกนาทีหรือห่างกว่านั้น มีวิธีที่สะอาดกว่า: เขียน activity ให้ fail เมื่อยังไม่เสร็จ แล้วให้ retry policy เป็นคนกำหนดจังหวะ เพราะ retry ของ activity ไม่ถูกบันทึกทีละครั้งใน history
func (a *Activities) CheckSettledOnce(ctx context.Context, ref BankTransferRef) (string, error) {
status, err := a.Bank.GetStatus(ctx, ref)
if err != nil {
return "", err
}
if status == "pending" {
// fail เพื่อให้ retry — นี่คือกลไก polling
return "", fmt.Errorf("still pending")
}
return status, nil
}
actCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: 30 * time.Second,
ScheduleToCloseTimeout: 72 * time.Hour, // เพดานรวมของการ poll
RetryPolicy: &temporal.RetryPolicy{
InitialInterval: time.Minute, // = ช่วงเวลา poll
BackoffCoefficient: 1.0, // 1.0 = ไม่ถอยห่าง คงที่ทุกนาที
MaximumAttempts: 0, // ไม่จำกัด ให้ ScheduleToClose เป็นตัวคุม
},
})
| Poll ถี่ (ลูปใน activity) | Poll ห่าง (retry เป็นตัวจับเวลา) | |
|---|---|---|
| ความถี่ที่เหมาะ | ทุกไม่กี่วินาที | ทุกนาทีขึ้นไป |
| ต้อง heartbeat | ต้อง | ไม่ต้อง |
| ทน worker ตาย | ต้องพึ่ง heartbeat + retry | ทนได้ดี |
| โค้ด | เขียนลูปเอง | เขียนแค่ตรวจครั้งเดียว |
Async Activity Completion
บางงานไม่ได้จบในตัว activity เอง เช่นส่งเรื่องเข้าระบบอนุมัติของอีกทีมแล้วรอเขากดกลับมา ซึ่งอาจเป็นวันถัดไป การถือ activity ค้างไว้ทั้งวันเปลืองเปล่าๆ
// Step 1: activity ส่งเรื่องแล้วบอกว่า "ผลจะมาทีหลัง"
func (a *Activities) RequestManualApproval(ctx context.Context, req WithdrawRequest) (string, error) {
token := activity.GetInfo(ctx).TaskToken
// เก็บ token ไว้ให้ระบบอื่นเอาไปใช้ตอบกลับ
if err := a.Approvals.Create(ctx, req.TransactionID, token); err != nil {
return "", err
}
return "", activity.ErrResultPending // ไม่ใช่ error — เป็นสัญญาณว่ารอผลภายนอก
}
// Step 2: อีกกระบวนการหนึ่ง (เช่น HTTP handler ตอน ops กดอนุมัติ) ปิดงานให้
func (h *Handler) OnApproved(w http.ResponseWriter, r *http.Request) {
token, err := h.Approvals.Token(r.Context(), txnID)
if err != nil {
http.Error(w, "not found", http.StatusNotFound)
return
}
if err := h.Temporal.CompleteActivity(r.Context(), token, ApprovalResult{
Approved: true,
ApproverID: currentUser(r),
}, nil); err != nil {
http.Error(w, "failed", http.StatusInternalServerError)
return
}
}
ปฏิเสธก็ส่ง error เข้าไปแทน:
_ = h.Temporal.CompleteActivity(ctx, token, nil,
temporal.NewNonRetryableApplicationError("rejected by compliance", "Rejected", nil))
ถ้าไม่อยากเก็บ task token ก็ระบุด้วย ID ได้:
_ = h.Temporal.CompleteActivityByID(ctx, namespace, workflowID, runID, activityID, result, nil)
เทียบกับ Signal แล้วเลือกให้ถูก
ถ้าระบบภายนอกเชื่อถือได้ว่าจะตอบกลับแน่ๆ และไม่ต้องการ heartbeat หรือ cancel ใช้ Signal ง่ายกว่ามาก (บทที่ 7) ไม่ต้องเก็บ task token
Async Activity Completion คุ้มเมื่อ: อยากได้ timeout ของ activity มาคุมโดยอัตโนมัติ, อยากให้ retry policy ทำงาน, หรืออยาก cancel งานที่ค้างได้จาก workflow
[!CAUTION] Task token ต้องเก็บอย่างปลอดภัย ใครถือ token นี้ = ปิดงานอนุมัติแทนได้ ต้องเก็บใน datastore ที่มีการควบคุมสิทธิ์ ไม่ใช่ส่งไปกับอีเมลหรือ query string และควรมีอายุจำกัด
Session — ปักงานหลายขั้นไว้กับ worker เครื่องเดียว
เฉพาะ Go SDK เท่านั้น ใช้เมื่อ activity หลายตัวต้องทำงานบนเครื่องเดียวกัน เพราะมี state อยู่บน disk ของเครื่องนั้น เช่นดาวน์โหลดไฟล์ statement มาประมวลผลแล้วอัปโหลดต่อ
w := worker.New(c, "recon-io", worker.Options{
EnableSessionWorker: true,
MaxConcurrentSessionExecutionSize: 50,
})
sessionCtx, err := workflow.CreateSession(ctx, &workflow.SessionOptions{
CreationTimeout: time.Minute,
ExecutionTimeout: 30 * time.Minute,
})
if err != nil {
return err
}
defer workflow.CompleteSession(sessionCtx)
// ทั้ง 3 ตัวรันบน worker เครื่องเดียวกันแน่นอน
var localPath string
if err := workflow.ExecuteActivity(sessionCtx, a.DownloadStatement, ref).
Get(sessionCtx, &localPath); err != nil {
return err
}
var summary ReconSummary
if err := workflow.ExecuteActivity(sessionCtx, a.ParseStatement, localPath).
Get(sessionCtx, &summary); err != nil {
return err
}
return workflow.ExecuteActivity(sessionCtx, a.UploadResult, summary).Get(sessionCtx, nil)
ข้อจำกัดของ Session
- ไม่รอดจาก worker restart — worker ตาย session ล้ม ได้
workflow.ErrSessionFailed - ไม่มีการรองรับฝั่ง server เลย Go SDK ทำเองทั้งหมดด้วย task queue ภายใน
- การจำกัดจำนวน session เป็นแบบ per-process ไม่ใช่ per-host
ก่อนใช้ Session ให้ถามก่อนว่ารวมเป็น activity ตัวเดียวที่มี heartbeat ได้ไหม ถ้าได้ ทำแบบนั้นง่ายกว่าและทนกว่า
Nexus — เรียกข้ามทีมข้าม namespace
Child workflow ใช้ได้เมื่ออยู่ใน namespace เดียวกัน แต่ถ้าทีม lending มี namespace ของตัวเองและเราอยากเรียกบริการของเขา Nexus คือทางที่ออกแบบมาให้
ฝั่งผู้ให้บริการประกาศ operation:
var creditService = nexus.NewService("credit")
var checkLimit = nexus.NewSyncOperation("check-limit",
func(ctx context.Context, req LimitRequest, _ nexus.StartOperationOptions) (LimitResult, error) {
return lookupLimit(ctx, req)
})
// operation ที่เบื้องหลังเป็น workflow ทั้งตัว
var approveLoan = temporalnexus.NewWorkflowRunOperation("approve-loan",
workflows.LoanApprovalWorkflow,
func(ctx context.Context, req LoanRequest, _ nexus.StartOperationOptions) (client.StartWorkflowOptions, error) {
return client.StartWorkflowOptions{ID: "loan-" + req.ApplicationID}, nil
})
_ = creditService.Register(checkLimit, approveLoan)
w.RegisterNexusService(creditService)
ฝั่งผู้เรียกใช้เหมือนเรียก activity:
nc := workflow.NewNexusClient("credit-endpoint", "credit")
var result LimitResult
if err := nc.ExecuteOperation(ctx, "check-limit", LimitRequest{WalletID: id},
workflow.NexusOperationOptions{}).Get(ctx, &result); err != nil {
return err
}
| Child Workflow | Nexus | |
|---|---|---|
| ขอบเขต | namespace เดียวกัน | ข้าม namespace / ข้ามทีม |
| สัญญาระหว่างกัน | ผูกกับ type ของ Go โดยตรง | ประกาศเป็น operation ชัดเจน |
| เหมาะกับ | งานย่อยของระบบเดียวกัน | เรียกบริการของทีมอื่นแบบมีขอบเขตชัด |
ประโยชน์เชิงองค์กรคือทีมผู้ให้บริการเปลี่ยน implementation ข้างในได้อิสระ ตราบใดที่ operation contract ยังเหมือนเดิม — ผู้เรียกไม่ต้องรู้ว่าข้างในเป็น workflow หรือไม่
Standalone Activities
ปกติ activity ต้องถูกเรียกจาก workflow เสมอ ฟีเจอร์นี้ให้เรียก activity ตรงๆ จาก client เหมาะกับงานที่อยากได้ retry และ visibility ของ Temporal แต่ไม่ต้องการ orchestration
handle, err := c.ExecuteActivity(ctx, client.ExecuteActivityOptions{
TaskQueue: "bank-io",
StartToCloseTimeout: 60 * time.Second,
}, activities.SyncExchangeRate, "THB")
if err != nil {
return err
}
var rate ExchangeRate
err = handle.Get(ctx, &rate)
Standalone Activities อยู่ในสถานะ Public Preview
ต้องใช้ Go SDK v1.41 ขึ้นไป, Temporal Server v1.31 ขึ้นไป และ CLI v1.7 ขึ้นไป ฟีเจอร์ที่ยังเป็น preview อาจเปลี่ยน API ได้ อย่าเพิ่งเอาไปวางบนเส้นทางที่เงินไหลผ่าน จนกว่าจะ GA ถ้าต้องใช้ตอนนี้ ให้จำกัดไว้กับงานเสริมอย่างการ sync อัตราแลกเปลี่ยน
สรุปว่าเมื่อไหร่ใช้อะไร
| โจทย์ | ใช้ |
|---|---|
| lookup เร็วๆ ที่ทำซ้ำแล้วไม่เสียหาย | Local Activity |
| ตรวจสถานะทุกไม่กี่วินาที | ลูป + heartbeat ใน activity |
| ตรวจสถานะทุกนาทีขึ้นไป | Retry policy BackoffCoefficient: 1.0 |
| รอคนกดอนุมัติ (เชื่อถือได้ว่าจะกด) | Signal |
| รอระบบภายนอกตอบ พร้อม timeout และ cancel | Async Activity Completion |
| หลาย activity ต้องอยู่เครื่องเดียวกัน | Session (หรือรวมเป็น activity เดียว) |
| เรียกบริการของทีมอื่น | Nexus |
| งานเดี่ยวที่อยากได้ retry ของ Temporal | Standalone Activity (ระวัง preview) |
ข้อควรระวังร่วมกันของทั้งบท
pattern ในบทนี้แลกความทนทานบางส่วนกับความเร็วหรือความยืดหยุ่นเกือบทุกตัว ก่อนหยิบมาใช้ ให้ถามเสมอว่า "ถ้าส่วนนี้ล้มกลางคัน เงินของลูกค้าจะอยู่ตรงไหน" ถ้าตอบไม่ได้ชัดเจน ให้กลับไปใช้ activity ปกติกับ saga จากบทที่ 6