บทที่ 8 · Part 2 — Building
Long-Running Workflows
Timer, Selector, Child Workflow, Continue-As-New, Entity Workflow และ Schedule สำหรับ reconciliation
ธุรกรรมส่วนใหญ่จบใน 30 วินาที แต่ระบบ wallet จริงมีงานที่กินเวลาเป็นวัน เป็นเดือน หรือไม่จบเลย — รอ settlement T+1, กระทบยอดตี 2 ทุกคืน, บัญชีลูกค้า 1 ใบที่มีชีวิตอยู่ตลอด บทนี้ว่าด้วยเครื่องมือที่ทำให้ workflow อยู่ได้นานโดยไม่ระเบิด
จบบทนี้คุณจะ
- ใช้ durable timer รอเป็นวันได้โดยไม่กิน resource
- ใช้ Selector จัดการเหตุการณ์ที่แข่งกันเข้ามาอย่างถูกต้อง
- รู้ว่าเมื่อไหร่ควรแตก Child Workflow และเมื่อไหร่แค่เขียนฟังก์ชันก็พอ
- ใช้ Continue-As-New โดยไม่ทำ signal หาย
- ออกแบบ Entity Workflow สำหรับบัญชี wallet และรู้ข้อจำกัดของมัน
- ตั้ง Schedule ให้ reconciliation รันทุกคืนพร้อม overlap policy ที่ถูก
Timer ที่รอดจาก worker restart
time.Sleep ใน workflow คือ non-determinism ตามที่ว่าไว้ในบทที่ 4
มันบล็อก goroutine จริงและหายไปพร้อม worker ที่ตาย
workflow.Sleep ไม่ใช่การหลับ — มันคือคำสั่งให้ server ตั้งนาฬิกาไว้ให้
worker ปล่อย workflow ออกจาก memory ได้เลย พอถึงเวลา server จะ schedule task ใหม่
แล้ว replay กลับมาที่บรรทัดถัดไป ระหว่างที่รอไม่มี process ไหนของเราค้างอยู่เลย
// รอครบ T+1 แล้วค่อยตรวจว่าเงินเข้าบัญชีปลายทางจริงหรือยัง
func SettlementCheckWorkflow(ctx workflow.Context, ref BankTransferRef) error {
if err := workflow.Sleep(ctx, 24*time.Hour); err != nil {
return err // ctx ถูก cancel ระหว่างรอ
}
actCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: 60 * time.Second,
})
var a *activities.Activities
var settled bool
if err := workflow.ExecuteActivity(actCtx, a.CheckSettlement, ref).Get(actCtx, &settled); err != nil {
return err
}
if !settled {
return workflow.ExecuteActivity(actCtx, a.RaiseOpsAlert, ref.TransactionID,
"not settled after T+1").Get(actCtx, nil)
}
return nil
}
รอเป็นเดือนก็ได้ ค่าใช้จ่ายเท่ากับ event 2 ตัวใน history (TimerStarted, TimerFired)
เลือกให้ถูกแบบ:
| ใช้ | เมื่อ | คืนค่า |
|---|---|---|
workflow.Sleep(ctx, d) | รอเฉยๆ ไม่มีอะไรมาขัด | error (non-nil ถ้าถูก cancel) |
workflow.NewTimer(ctx, d) | ต้องเอาไปใส่ Selector แข่งกับอย่างอื่น | workflow.Future |
workflow.AwaitWithTimeout(ctx, d, cond) | รอจนเงื่อนไขเป็นจริง หรือหมดเวลา | (bool, error) — false แปลว่าหมดเวลา |
// รอ ops อนุมัติภายใน 8 ชั่วโมง — สั้นกว่าเขียน Selector เองเมื่อเงื่อนไขเป็น bool ธรรมดา
approved := false
workflow.Go(ctx, func(gctx workflow.Context) {
workflow.GetSignalChannel(gctx, "ops-approved").Receive(gctx, &approved)
})
ok, err := workflow.AwaitWithTimeout(ctx, 8*time.Hour, func() bool { return approved })
if err != nil {
return err
}
if !ok {
return temporal.NewNonRetryableApplicationError(
"approval timed out", "ApprovalTimeout", nil)
}
Timer ไม่ใช่นาฬิกาที่แม่นระดับวินาที
Temporal รับประกันว่า timer จะ ไม่ยิงก่อนเวลา แต่ยิงช้ากว่ากำหนดได้ ตามภาระของ cluster และจังหวะที่ worker ว่าง อย่าใช้ timer เป็นตัวกำหนด cutoff ทางบัญชีแบบเป๊ะๆ (เช่น "ต้องตัดยอด 23:59:59.000") ให้ activity อ่านเวลาจริงจากระบบบัญชีแล้วตัดสินเองแทน
[!TIP] เวลาใน workflow ต้องมาจาก workflow.Now(ctx)
ค่านี้ถูกบันทึกใน history จึงได้ค่าเดิมทุกครั้งที่ replay
time.Now() จะเปลี่ยนไปทุกรอบและทำให้ workflow พัง
Selector — จัดการเหตุการณ์ที่แข่งกัน
บทที่ 7 ใช้ Selector รอ "OTP หรือ timeout อันไหนมาก่อน" มาแล้ว ตรงนี้ขยายไปที่กรณีซับซ้อนขึ้น: รอ callback จากธนาคาร แต่ระหว่างรอก็ poll สถานะเป็นระยะ และยอมให้ ops สั่งยกเลิกได้ตลอด
func AwaitBankResultWorkflow(ctx workflow.Context, ref BankTransferRef) (string, error) {
actCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: 30 * time.Second,
})
var a *activities.Activities
callbackCh := workflow.GetSignalChannel(ctx, "bank-callback")
abortCh := workflow.GetSignalChannel(ctx, "ops-abort")
deadline := workflow.NewTimer(ctx, 6*time.Hour)
outcome := ""
for outcome == "" {
// poll timer ต้องสร้างใหม่ทุกรอบ — Future 1 ตัวยิงได้ครั้งเดียว
poll := workflow.NewTimer(ctx, 5*time.Minute)
sel := workflow.NewSelector(ctx)
sel.AddReceive(callbackCh, func(c workflow.ReceiveChannel, _ bool) {
var cb BankCallback
c.Receive(ctx, &cb)
outcome = cb.Status
})
sel.AddReceive(abortCh, func(c workflow.ReceiveChannel, _ bool) {
c.Receive(ctx, nil)
outcome = "aborted"
})
sel.AddFuture(poll, func(f workflow.Future) {
if f.Get(ctx, nil) != nil {
return // timer ถูก cancel
}
var status string
if err := workflow.ExecuteActivity(actCtx, a.PollBankStatus, ref).Get(actCtx, &status); err == nil &&
status != "pending" {
outcome = status
}
})
sel.AddFuture(deadline, func(f workflow.Future) {
if f.Get(ctx, nil) == nil {
outcome = "timeout"
}
})
sel.Select(ctx)
}
return outcome, nil
}
จุดที่พลาดกันบ่อย:
| พฤติกรรม | ผลที่ตามมา |
|---|---|
AddFuture ยิงได้ ครั้งเดียว ต่อ Future 1 ตัว | timer ที่ต้องเต้นซ้ำ ต้อง NewTimer ใหม่ทุกรอบลูป |
Select ประมวลผลทีละ 1 branch | ถ้ามีหลายเหตุการณ์พร้อมกัน ต้องวนเรียก Select ซ้ำ |
callback ของ AddReceive ต้อง เรียก c.Receive | ถ้าไม่รับ ข้อความยังค้างอยู่ในคิว แล้วจะยิงซ้ำไม่รู้จบ |
AddDefault ยิงทันทีเมื่อไม่มี branch ไหนพร้อม | ทำให้ Select ไม่บล็อก — ใช้ตรวจ "มีอะไรค้างไหม" ไม่ใช่วนลูปเปล่า |
sel.HasPending() | เช็คว่ายังมีข้อความค้างก่อนจะจบ workflow |
ห้ามใช้ `select` ของ Go ใน workflow
select, time.After, sync.WaitGroup, channel ธรรมดา และ go func()
ทั้งหมดนี้ไม่ผ่าน replay ต้องใช้คู่แฝดของ Temporal เสมอ:
workflow.Selector, workflow.NewTimer, workflow.NewChannel, workflow.Go
Child Workflow
Child Workflow คือ workflow ที่ถูกสั่งจาก workflow อีกตัว มี history ของตัวเอง มี retry policy ของตัวเอง และล้มได้โดยไม่ลาก parent ไปด้วย
อย่าแตก child เพียงเพราะอยากแยกโค้ด
การจัดระเบียบโค้ดใช้ฟังก์ชันธรรมดาก็พอ — เรียก chargeBank(ctx, req) ที่เป็น
ฟังก์ชัน Go ปกติได้เลย ไม่มีค่าใช้จ่ายอะไรเพิ่ม
child workflow มีต้นทุน (event เพิ่ม, latency เพิ่ม, ID ต้องจัดการ)
ใช้เมื่อมีเหตุผลเชิงระบบเท่านั้น
ใช้เมื่อ:
- กัน history บวม — parent สั่งงาน 500 รายการ ถ้าทำในตัวเองจะได้ event หลายพัน
- แยก failure domain — ธนาคาร A ล่ม ไม่ควรทำให้รอบกระทบยอดของธนาคาร B พัง
- คนละ retry / timeout policy — งาน batch ยอมช้าได้ งาน real-time ยอมไม่ได้
- อายุไม่เท่ากัน — parent จบแล้วแต่ลูกต้องอยู่ต่อ (
ABANDON)
เคสจริง: รอบกระทบยอดรายวัน แตกเป็น 1 child ต่อ 1 ธนาคาร
func DailyReconciliationWorkflow(ctx workflow.Context, day string) (ReconResult, error) {
banks := []string{"SCB", "KBANK", "BBL"}
// ยิงพร้อมกันทั้งหมดก่อน แล้วค่อยเก็บผล — ไม่งั้นจะกลายเป็นทำทีละธนาคาร
futures := make(map[string]workflow.ChildWorkflowFuture, len(banks))
for _, bank := range banks {
bankCtx := workflow.WithChildOptions(ctx, workflow.ChildWorkflowOptions{
// ตั้ง ID เองเสมอ เพื่อให้ ops ค้นเจอใน Web UI และกันรันซ้ำรอบเดิม
WorkflowID: fmt.Sprintf("recon-%s-%s", day, bank),
TaskQueue: "recon",
WorkflowExecutionTimeout: 4 * time.Hour,
RetryPolicy: &temporal.RetryPolicy{
MaximumAttempts: 3,
},
})
futures[bank] = workflow.ExecuteChildWorkflow(bankCtx, ReconcileBankWorkflow, day, bank)
}
result := ReconResult{Day: day, PerBank: map[string]BankReconResult{}}
for bank, f := range futures {
var r BankReconResult
if err := f.Get(ctx, &r); err != nil {
// ธนาคารเดียวพัง ไม่ล้มทั้งรอบ — บันทึกไว้แล้วไปต่อ
workflow.GetLogger(ctx).Error("bank recon failed", "bank", bank, "error", err)
result.Failed = append(result.Failed, bank)
continue
}
result.PerBank[bank] = r
result.MismatchSatang += r.MismatchSatang
}
return result, nil
}
Options ที่ต้องรู้
| Option | ทำอะไร |
|---|---|
WorkflowID | ตั้งเองเสมอ — deterministic และค้นหาได้ ถ้าไม่ตั้ง SDK สร้างให้แบบสุ่ม(แต่ replay-safe) |
TaskQueue | ว่างไว้ = ใช้ queue เดียวกับ parent ตั้งได้ถ้าอยากแยก worker pool |
WorkflowExecutionTimeout | เวลารวมทั้งหมดของ child รวม retry |
RetryPolicy | child ไม่ retry โดย default ต่างจาก activity |
ParentClosePolicy | ชะตากรรมของ child เมื่อ parent ปิดตัว |
WorkflowIDReusePolicy | กันรันซ้ำ ID เดิม — สำคัญมากกับ ID ที่อิงวันที่ |
WaitForCancellation | ถ้า parent ถูก cancel ให้รอ child เก็บกวาดเสร็จก่อนไหม |
ParentClosePolicy มี 3 ค่า:
| ค่า | ผล | ใช้กับ |
|---|---|---|
PARENT_CLOSE_POLICY_TERMINATE (default) | child ถูก terminate ทันที | งานที่ไม่มีความหมายถ้า parent จบแล้ว |
PARENT_CLOSE_POLICY_REQUEST_CANCEL | ส่ง cancel ให้ child จัดการเอง | เลือกตัวนี้ถ้า child แตะเงิน — จะได้ compensate ทัน |
PARENT_CLOSE_POLICY_ABANDON | child อยู่ต่ออิสระ | งานยาวที่ต้องจบเองเช่นการติดตาม settlement |
`TERMINATE` = ตายทันที ไม่มี compensation
เหมือน temporal workflow terminate ที่พูดถึงในบทที่ 6
child ที่กำลังจะโอนเงินอยู่แล้วโดน terminate จะทิ้งสถานะค้างไว้กลางทาง
ถ้า child แตะเงิน ให้ใช้ REQUEST_CANCEL เสมอ
ยิงแล้วไม่รอ
ถ้าอยากให้ child อยู่ต่อหลัง parent จบ ต้องรอให้มัน เริ่ม จริงก่อน ไม่งั้น parent อาจจบก่อนที่คำสั่งสร้าง child จะถูกส่งออกไป
detachedCtx := workflow.WithChildOptions(ctx, workflow.ChildWorkflowOptions{
WorkflowID: "settlement-check-" + ref.TransactionID,
ParentClosePolicy: enumspb.PARENT_CLOSE_POLICY_ABANDON,
})
child := workflow.ExecuteChildWorkflow(detachedCtx, SettlementCheckWorkflow, ref)
// รอแค่ "เริ่มแล้ว" ไม่ต้องรอผลลัพธ์
var exec workflow.Execution
if err := child.GetChildWorkflowExecution().Get(ctx, &exec); err != nil {
return err
}
workflow.GetLogger(ctx).Info("settlement check started", "runID", exec.RunID)
ส่ง signal เข้า child ได้ตรงๆ ผ่าน future เดียวกัน:
err := child.SignalChildWorkflow(ctx, "bank-callback", callbackPayload).Get(ctx, nil)
Continue-As-New
Event History โตขึ้นเรื่อยๆ ตามจำนวนสิ่งที่ workflow ทำ และ server มีเพดานทั้งจำนวน event และขนาดรวม พอชนเพดาน workflow จะถูกบังคับให้ล้ม — สำหรับ workflow ที่วนลูปไม่รู้จบ นี่ไม่ใช่ "ถ้าเกิด" แต่คือ "เมื่อไหร่"
Continue-As-New คือการปิด run ปัจจุบันแล้วเปิด run ใหม่ทันที ด้วย Workflow ID เดิม แต่ history ว่างเปล่า สถานะที่อยากเก็บต้องส่งต่อไปเป็น input
type WalletState struct {
WalletID string
BalanceSatang int64
// เก็บเท่าที่จำเป็นต่อการทำงานรอบถัดไป — ห้ามสะสมประวัติทั้งหมดไว้ตรงนี้
LastTxnID string
DailyUsed int64
DailyResetAt time.Time
}
func WalletAccountWorkflow(ctx workflow.Context, state WalletState) error {
// ... ทำงานเป็นลูป ...
if workflow.GetInfo(ctx).GetContinueAsNewSuggested() {
return workflow.NewContinueAsNewError(ctx, WalletAccountWorkflow, state)
}
return nil
}
GetContinueAsNewSuggested() ให้ server เป็นคนบอกว่าใกล้เพดานแล้ว ซึ่งดีกว่า hardcode ตัวเลข
เพราะเพดานเป็นค่า config ระดับ cluster ที่เปลี่ยนได้ ถ้าต้องการเกณฑ์ของตัวเองเพิ่ม
(เช่นอยากตัดถี่กว่านั้นเพื่อให้ replay เร็ว) ใช้ GetCurrentHistoryLength() ประกอบ:
func shouldContinueAsNew(ctx workflow.Context) bool {
info := workflow.GetInfo(ctx)
if info.GetContinueAsNewSuggested() {
return true
}
return info.GetCurrentHistoryLength() > 2_000
}
3 กับดักของ Continue-As-New
- Signal หายได้ — signal ที่ค้างอยู่ในคิวตอน continue จะไม่ถูกส่งต่อไป run ใหม่ ต้อง drain ให้หมดก่อน
- State ไม่ถูกส่งต่อเอง — ตัวแปรทั้งหมดหายไป เหลือแค่สิ่งที่ใส่ใน input
- Input โตไม่จำกัดไม่ได้ — ถ้ายัด transaction ทั้งหมดที่เคยทำลงใน state เท่ากับย้ายปัญหา history บวมไปเป็น payload บวมแทน (ดูเรื่อง payload limit ในบทที่ 11)
Drain ให้ถูก:
func drainAndContinue(ctx workflow.Context, state WalletState, ch workflow.ReceiveChannel) error {
// 1. รอ signal/update handler ที่ยังทำงานค้างให้เสร็จก่อน (TMPRL1102)
if err := workflow.Await(ctx, func() bool {
return workflow.AllHandlersFinished(ctx)
}); err != nil {
return err
}
// 2. ดูดข้อความที่ค้างในคิวออกให้หมด แล้วประมวลผลตามปกติ
for {
var cmd WalletCommand
if !ch.ReceiveAsync(&cmd) {
break
}
state = applyCommand(state, cmd)
}
// 3. ค่อย continue
return workflow.NewContinueAsNewError(ctx, WalletAccountWorkflow, state)
}
ห้ามเรียก Continue-As-New จากใน signal/update handler
ต้องให้ handler ตั้ง flag แล้วให้ลูปหลักเป็นคนตัดสินใจ continue
ถ้า handler เป็นคน return NewContinueAsNewError handler ตัวอื่นที่ทำงานอยู่จะถูกตัดกลางคัน
Entity Workflow
Entity Workflow คือการมองว่า "1 บัญชี = 1 workflow ที่ไม่จบ" แทนที่จะเปิด workflow ใหม่ทุกครั้งที่มีธุรกรรม เราเปิดตัวเดียวต่อ wallet แล้วยิงคำสั่งเข้าไปด้วย signal/update ตลอดอายุการใช้งาน
ประโยชน์ในบริบทการเงินคือได้ serialization ฟรี — คำสั่งทั้งหมดของ wallet ใบนั้น ถูกประมวลผลทีละคำสั่งตามลำดับใน workflow เดียว ไม่มี race, ไม่ต้องล็อก DB, และตรวจ daily limit ได้จาก state ในหน่วยความจำโดยตรง
func WalletAccountWorkflow(ctx workflow.Context, state WalletState) error {
logger := workflow.GetLogger(ctx)
// อ่านสถานะได้ตลอดจากหน้าจอ ops
if err := workflow.SetQueryHandler(ctx, "balance", func() (int64, error) {
return state.BalanceSatang, nil
}); err != nil {
return err
}
// สั่งเติมเงินแล้วรู้ผลทันที — คู่กับ Update-with-Start จากบทที่ 7
err := workflow.SetUpdateHandlerWithOptions(ctx, "top-up",
func(ctx workflow.Context, cmd TopUpCommand) (TopUpAccepted, error) {
// idempotency: คำสั่งเดิมส่งซ้ำต้องได้ผลเดิม ไม่ใช่เติม 2 รอบ
if cmd.TransactionID == state.LastTxnID {
return TopUpAccepted{BalanceSatang: state.BalanceSatang, Duplicate: true}, nil
}
actCtx := workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: 60 * time.Second,
})
var a *activities.Activities
var bankRef string
if err := workflow.ExecuteActivity(actCtx, a.ChargeBank, cmd).Get(actCtx, &bankRef); err != nil {
return TopUpAccepted{}, err
}
state.BalanceSatang += cmd.AmountSatang
state.DailyUsed += cmd.AmountSatang
state.LastTxnID = cmd.TransactionID
return TopUpAccepted{BalanceSatang: state.BalanceSatang, BankRef: bankRef}, nil
},
workflow.UpdateHandlerOptions{
Validator: func(ctx workflow.Context, cmd TopUpCommand) error {
if cmd.AmountSatang <= 0 {
return fmt.Errorf("amount must be positive")
}
if state.DailyUsed+cmd.AmountSatang > 100_000_00 {
return fmt.Errorf("daily limit exceeded")
}
return nil
},
},
)
if err != nil {
return err
}
closeCh := workflow.GetSignalChannel(ctx, "close-account")
closed := false
workflow.Go(ctx, func(gctx workflow.Context) {
closeCh.Receive(gctx, nil)
closed = true
})
// รีเซ็ตยอดใช้รายวัน และตรวจว่าถึงเวลา continue-as-new หรือยัง
for {
nextReset := workflow.Now(ctx).Truncate(24 * time.Hour).Add(24 * time.Hour)
fired, err := workflow.AwaitWithTimeout(ctx, nextReset.Sub(workflow.Now(ctx)),
func() bool { return closed || workflow.GetInfo(ctx).GetContinueAsNewSuggested() })
if err != nil {
return err
}
if !fired {
state.DailyUsed = 0
state.DailyResetAt = workflow.Now(ctx)
continue
}
if closed {
logger.Info("wallet closed", "walletID", state.WalletID)
return workflow.Await(ctx, func() bool { return workflow.AllHandlersFinished(ctx) })
}
if err := workflow.Await(ctx, func() bool {
return workflow.AllHandlersFinished(ctx)
}); err != nil {
return err
}
return workflow.NewContinueAsNewError(ctx, WalletAccountWorkflow, state)
}
}
ฝั่ง API ใช้ Update-with-Start เพื่อ "เปิดถ้ายังไม่มี แล้วสั่งเลย" ตามที่อธิบายไว้ใน บทที่ 7 จึงไม่ต้องมีขั้นตอน provision บัญชีแยกต่างหาก
ข้อจำกัดที่ต้องรู้ก่อนเลือก Entity Workflow
- Throughput ต่อ workflow มีเพดาน — คำสั่งของ wallet ใบเดียวถูกประมวลผลแบบ serialize ตามลำดับ ประมาณหลักสิบครั้งต่อวินาทีเป็นตัวเลขที่ควรออกแบบเผื่อไว้ บัญชีร้านค้าที่มีธุรกรรม 1,000 ครั้งต่อวินาทีไม่เหมาะกับรูปแบบนี้
- จำนวน workflow ที่ Running พร้อมกันจะเยอะมาก — 1 ใบต่อลูกค้า 1 คน ต้องคุยกับทีม infra เรื่อง capacity ก่อน
- การแก้โค้ดยากขึ้น — workflow ที่รันมา 6 เดือนต้อง replay ผ่านโค้ดใหม่ได้ เรื่องนี้คือหัวใจของบทที่ 10
ทางเลือกที่ปลอดภัยกว่าและใช้กันมากคือ workflow ต่อ 1 ธุรกรรม (แบบบทที่ 3–6) แล้วเก็บ balance ไว้ใน DB ตามเดิม ใช้ entity workflow เฉพาะที่ต้องการ serialization จริงๆ เท่านั้น
Schedule — งานที่ต้องรันตามเวลา
อย่าเขียน for { workflow.Sleep(24h); doWork() } เพื่อทำงานประจำวัน
มันดูง่ายแต่ไม่มี UI ให้ ops กด pause, ย้อนหลังไม่ได้, และถ้ารอบก่อนยังไม่จบก็ไม่มีใครห้าม
Schedule เป็น object ระดับ server ที่จัดการเรื่องพวกนี้ให้:
handle, err := c.ScheduleClient().Create(context.Background(), client.ScheduleOptions{
ID: "daily-reconciliation",
Spec: client.ScheduleSpec{
Calendars: []client.ScheduleCalendarSpec{
{
Hour: []client.ScheduleRange{{Start: 2}},
Minute: []client.ScheduleRange{{Start: 0}},
Comment: "กระทบยอดรายวันหลังปิดรอบธนาคาร",
},
},
TimeZoneName: "Asia/Bangkok", // สำคัญ — default คือ UTC
Jitter: 5 * time.Minute,
},
Action: &client.ScheduleWorkflowAction{
ID: "recon", // จะถูกต่อท้ายด้วย timestamp ของแต่ละรอบ
Workflow: workflows.DailyReconciliationWorkflow,
TaskQueue: "recon",
WorkflowExecutionTimeout: 4 * time.Hour,
},
Overlap: enumspb.SCHEDULE_OVERLAP_POLICY_SKIP,
CatchupWindow: 2 * time.Hour,
PauseOnFailure: true,
})
ตั้ง `TimeZoneName` เสมอ
ค่า default เป็น UTC ถ้าลืม รอบ "ตีสอง" จะไปรันตอน 9 โมงเช้าเวลาไทย
และอย่าลืมว่า Jitter ช่วยไม่ให้ทุก schedule ยิงพร้อมกันจนถล่ม worker pool
Overlap Policy — ข้อที่สำคัญที่สุดสำหรับงานการเงิน
ถ้ารอบเมื่อวานยังไม่จบแล้วถึงเวลารอบใหม่ จะทำยังไง
| Policy | พฤติกรรม | เหมาะกับ |
|---|---|---|
SKIP (default) | ข้ามรอบใหม่ไปเลย | งานกระทบยอด — รันซ้อนกันแล้วตัวเลขจะมั่ว |
BUFFER_ONE | เก็บรอบใหม่ไว้ 1 รอบ รันต่อเมื่อรอบเดิมจบ | งานที่ห้ามข้ามแต่รันพร้อมกันไม่ได้ |
BUFFER_ALL | ต่อคิวทุกรอบที่ค้าง | งาน batch ที่ทุกรอบต้องได้รัน |
CANCEL_OTHER | ยกเลิกรอบเดิมแล้วเริ่มรอบใหม่ | งานที่ข้อมูลล่าสุดสำคัญกว่างานที่ค้าง |
TERMINATE_OTHER | ฆ่ารอบเดิมทันที | ระวัง — ไม่มี compensation |
ALLOW_ALL | รันซ้อนได้ไม่จำกัด | งาน read-only ล้วน |
PauseOnFailure: true ทำให้ schedule หยุดตัวเองเมื่อรอบหนึ่งล้ม
คุ้มค่ามากกับงานกระทบยอด เพราะรอบถัดไปที่รันบนฐานข้อมูลที่ยังไม่ถูกแก้ก็จะล้มซ้ำอยู่ดี
แถมยังกลบร่องรอยของรอบแรก
จัดการ schedule
h := c.ScheduleClient().GetHandle(context.Background(), "daily-reconciliation")
_ = h.Pause(ctx, client.SchedulePauseOptions{Note: "ปิดระบบบำรุงรักษา"})
_ = h.Unpause(ctx, client.ScheduleUnpauseOptions{Note: "เสร็จแล้ว"})
// สั่งรันเดี๋ยวนี้เลย โดยไม่รบกวนตารางปกติ
_ = h.Trigger(ctx, client.ScheduleTriggerOptions{})
// ย้อนรันรอบที่ขาดไป เช่นตอน worker ล่มไป 3 วัน
_ = h.Backfill(ctx, client.ScheduleBackfillOptions{
Backfill: []client.ScheduleBackfill{{
Start: time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC),
End: time.Date(2026, 8, 3, 0, 0, 0, 0, time.UTC),
Overlap: enumspb.SCHEDULE_OVERLAP_POLICY_BUFFER_ALL,
}},
})
desc, _ := h.Describe(ctx)
จาก CLI:
temporal schedule create \
--schedule-id daily-reconciliation \
--calendar '{"hour":"2","minute":"0"}' \
--workflow-id recon \
--task-queue recon \
--type DailyReconciliationWorkflow \
--overlap-policy Skip
temporal schedule describe --schedule-id daily-reconciliation
temporal schedule toggle --schedule-id daily-reconciliation --pause --reason "maintenance"
temporal schedule trigger --schedule-id daily-reconciliation
temporal schedule backfill --schedule-id daily-reconciliation \
--start-time 2026-08-01T00:00:00Z \
--end-time 2026-08-03T00:00:00Z \
--overlap-policy BufferAll
Schedule กับ Cron Workflow ไม่ใช่ตัวเดียวกัน
StartWorkflowOptions.CronSchedule เป็นของเก่าที่ยังใช้ได้แต่ไม่แนะนำสำหรับของใหม่
Schedule ให้ pause/backfill/trigger/overlap policy และมีหน้า UI ของตัวเอง
ซึ่ง cron ไม่มี ถ้าเริ่มโปรเจกต์ใหม่ให้ใช้ Schedule
เลือกเครื่องมือให้ถูกงาน
| โจทย์ | ใช้ |
|---|---|
| รอเวลาที่แน่นอน | workflow.Sleep |
| รอเงื่อนไข หรือหมดเวลา | workflow.AwaitWithTimeout |
| รอหลายเหตุการณ์ที่แข่งกัน | workflow.Selector + workflow.NewTimer |
| งานย่อยจำนวนมากที่ล้มแยกกันได้ | Child Workflow |
| ลูปไม่รู้จบ / history โต | Continue-As-New |
| สถานะที่มีชีวิตยาว ต้อง serialize คำสั่ง | Entity Workflow (+ Continue-As-New) |
| งานตามเวลา | Schedule |
Checklist ก่อนปล่อย workflow ที่รันยาว
- ทุก loop ที่ไม่มีเงื่อนไขจบชัดเจน มีการเช็ค
GetContinueAsNewSuggested() - ก่อน Continue-As-New มี
AllHandlersFinishedและ drain signal channel ครบทุกช่อง - State ที่ส่งข้าม run มีขนาดคงที่ ไม่โตตามจำนวนธุรกรรม
- Child ที่แตะเงินใช้
ParentClosePolicyเป็นREQUEST_CANCELไม่ใช่TERMINATE - Child Workflow ID เป็น deterministic และค้นหาได้จาก Web UI
- Schedule ตั้ง
TimeZoneNameและOverlapตรงกับความหมายทางบัญชีของงานนั้น - Timer ที่ไม่ต้องใช้แล้วถูก cancel
ทั้งหมดนี้ทดสอบได้จริงด้วย time-skipping test environment — เรื่องของบทที่ 9