บทที่ 5 · Part 2 — Building
Activities in Depth
Timeout ทั้ง 4 ชนิด, Retry Policy, Heartbeat และ Idempotency กับ bank API ที่ตอบ timeout
Activity คือจุดที่เงินเคลื่อนที่จริง การตั้งค่า timeout กับ retry ผิด ไม่ได้แค่ทำให้ช้า แต่ทำให้ลูกค้าถูกหักเงินซ้ำได้
จบบทนี้คุณจะ
- เลือก timeout ทั้ง 4 ชนิดได้ถูกต้อง
- ออกแบบ Retry Policy ที่เหมาะกับ bank API
- ทำ activity ให้ idempotent ได้จริง ไม่ใช่แค่บนกระดาษ
- ใช้ heartbeat กับงานยาว และรองรับการ cancel
Timeout ทั้ง 4 ชนิด
HeartbeatTimeout เป็นอีกมิติหนึ่ง คือระยะห่างสูงสุดที่ยอมให้เกิดขึ้นระหว่าง heartbeat 2 ครั้ง
| Timeout | ควบคุมอะไร | คำแนะนำ |
|---|---|---|
StartToCloseTimeout | เวลาสูงสุดของการพยายาม1 ครั้ง | ตัวหลักที่ควรตั้งเสมอ ตั้งให้ยาวกว่า p99 ของ API ปลายทางพอสมควร |
ScheduleToCloseTimeout | เวลารวมทั้งหมดรวม retry ทุกครั้ง | ใช้เมื่อมี business deadline จริง เช่น "ต้องจบก่อนปิด batch 23:00" |
ScheduleToStartTimeout | เวลารอในคิวก่อนมี worker มาหยิบ | ไม่ค่อยจำเป็น ถ้าตั้งแล้ว timeout แปลว่า worker ไม่พอ ซึ่ง retry ไม่ช่วยอะไร |
HeartbeatTimeout | ระยะห่างสูงสุดระหว่าง heartbeat 2 ครั้ง | จำเป็นสำหรับ activity ที่รันนานกว่าไม่กี่วินาที |
ต้องตั้งอย่างน้อยหนึ่งใน StartToCloseTimeout หรือ ScheduleToCloseTimeout ไม่งั้น SDK จะ error
กับดักคลาสสิกในระบบการเงิน
ตั้ง StartToCloseTimeout: 5 * time.Second สำหรับ activity ที่เรียก bank API
ธนาคารตอบช้า 8 วินาทีแต่ทำรายการสำเร็จแล้ว →
Temporal ถือว่า timeout → retry → หักเงินอีกรอบ
ป้องกันได้ 2 ชั้น: (1) ตั้ง timeout ให้ยาวกว่า timeout ฝั่งธนาคารเสมอ (2) ส่ง idempotency key ทุกครั้ง — ชั้นที่ 2 คือชั้นที่คุณเชื่อถือได้จริง
Retry Policy
ctx = workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: 30 * time.Second,
RetryPolicy: &temporal.RetryPolicy{
InitialInterval: time.Second,
BackoffCoefficient: 2.0, // 1s, 2s, 4s, 8s, ...
MaximumInterval: time.Minute,
MaximumAttempts: 0, // 0 = ไม่จำกัด
NonRetryableErrorTypes: []string{"LimitExceeded", "AccountFrozen"},
},
})
ค่าเริ่มต้นของ SDK คือ retry แบบไม่จำกัดด้วย exponential backoff ซึ่งเป็นค่าที่ดีสำหรับงานส่วนใหญ่
อย่าตั้งค่าเกินจำเป็น
ปรับ MaximumAttempts หรือ MaximumInterval ก็ต่อเมื่อมีเหตุผลทางธุรกิจรองรับ
การจำกัด attempt ไว้ต่ำๆ ทำให้ธุรกรรมล้มตอนธนาคารสะดุดแค่แป๊บเดียว
ซึ่งมักแย่กว่าการปล่อยให้รอ
เลือก policy ตามลักษณะงาน
| Activity | Policy ที่แนะนำ | เหตุผล |
|---|---|---|
| เรียก bank API (หักเงิน/โอน) | retry ไม่จำกัด + backoff + idempotency key | ธนาคารล่มเป็นเรื่องชั่วคราว ยอมรอดีกว่าล้มธุรกรรม |
| ตรวจ limit / KYC | retry ปกติ แต่ error เชิงกฎธุรกิจต้อง non-retryable | "เกิน limit" retry อีกล้านครั้งก็ไม่ผ่าน |
| เขียน ledger | retry ไม่จำกัด | DB ล่มคือปัญหาชั่วคราว และห้ามยอมให้ขั้นนี้ล้มเด็ดขาด |
| ส่ง push / SMS | MaximumAttempts: 3 | ไม่กระทบเงิน ล้มแล้วปล่อยผ่านได้ |
| เช็คสถานะจากธนาคาร (polling) | BackoffCoefficient: 1 + InitialInterval = ช่วงที่อยากถาม | ทำให้ retry กลายเป็น polling loop ที่ไม่กิน history — ดูบทที่ 13 |
แยก error ให้ถูกประเภท
นี่คือการตัดสินใจที่สำคัญที่สุดในการเขียน activity — ผิดแล้วเจ็บทั้ง 2ทาง
| Retryable | Non-retryable |
|---|---|
| network error, connection reset | input ไม่ถูกต้อง (validation) |
| timeout, 503, 502 | authentication / authorization ล้มเหลว |
| rate limit (429) | ยอดเงินไม่พอ |
| DB deadlock, connection pool หมด | บัญชีถูกอายัด / ปิด |
| ธนาคารแจ้ง "ระบบปิดปรับปรุง" | เกิน daily limit |
| ตอบกลับ 500 แบบไม่ระบุสาเหตุ | ธนาคารปฏิเสธรายการ (rejected) |
func (a *Activities) ChargeBank(ctx context.Context, req TopUpRequest) (string, error) {
resp, err := a.Bank.Debit(ctx, BankDebitRequest{
IdempotencyKey: req.TransactionID,
AccountID: req.BankAccountID,
AmountSatang: req.AmountSatang,
})
if err != nil {
var apiErr *BankAPIError
if errors.As(err, &apiErr) {
switch apiErr.Code {
case "INSUFFICIENT_FUNDS", "ACCOUNT_CLOSED", "ACCOUNT_FROZEN":
// ธนาคารตัดสินใจแล้ว retry ไม่เปลี่ยนคำตอบ
return "", temporal.NewNonRetryableApplicationError(
apiErr.Message, "BankRejected", err,
)
case "RATE_LIMITED", "SERVICE_UNAVAILABLE":
// ปัญหาชั่วคราว — ปล่อยให้ Temporal retry
return "", err
}
}
return "", err // ไม่รู้จัก → ถือว่า retryable ปลอดภัยกว่า
}
return resp.Reference, nil
}
อีกทางเลือกคือประกาศไว้ที่ policy แทน ซึ่งทำให้เห็นภาพรวมจากฝั่ง workflow:
RetryPolicy: &temporal.RetryPolicy{
NonRetryableErrorTypes: []string{"BankRejected", "LimitExceeded"},
}
เมื่อไม่แน่ใจ ให้เลือก retryable
การ retry งานที่ idempotent เกินความจำเป็น = เสีย compute นิดหน่อย การไม่ retry งานที่จริงๆ แค่สะดุดชั่วคราว = ธุรกรรมลูกค้าล้มโดยไม่จำเป็น ต้นทุนไม่เท่ากันอย่างชัดเจน
Idempotency — หัวใจของงานการเงิน
สมมติฐานที่ต้องยึด
ทุก activity จะถูกเรียกมากกว่า 1 ครั้ง — จาก retry, จาก worker ตายหลังทำงานเสร็จแต่ก่อนรายงานผล, หรือจากการ reset ถ้าโค้ดของคุณทนกับสิ่งนี้ไม่ได้ แปลว่ามันมีบั๊กแล้ว
รูปแบบที่ 1 — Idempotency key ส่งไปให้ระบบปลายทาง
ดีที่สุดถ้าธนาคารรองรับ ต้องเป็นค่าที่ deterministic คือคำนวณได้เหมือนเดิมทุกครั้ง ไม่ใช่สุ่มใหม่
| แหล่งที่ดี | แหล่งที่ห้ามใช้ |
|---|---|
| Transaction ID จากระบบของคุณ | uuid.New() สร้างใหม่ใน activity |
workflow.GetInfo(ctx).WorkflowExecution.ID | time.Now() ประกอบเป็น key |
| Workflow ID + ชื่อขั้นตอน | activity.GetInfo(ctx).Attempt (เปลี่ยนทุก retry!) |
รูปแบบที่ 2 — Check before act
ใช้เมื่อ API ปลายทางไม่มี idempotency key ให้ (ซึ่งพบบ่อยกับ core banking รุ่นเก่า)
func (a *Activities) TransferToBank(ctx context.Context, req TransferRequest) (string, error) {
// 1) ถามก่อนว่าเคยทำรายการนี้ไปหรือยัง
existing, err := a.Bank.QueryByReference(ctx, req.TransactionID)
if err != nil && !errors.Is(err, ErrNotFound) {
return "", err
}
if existing != nil {
activity.GetLogger(ctx).Info("already transferred, reusing result",
"ref", existing.Reference)
return existing.Reference, nil
}
// 2) ยังไม่เคย → ทำจริง
resp, err := a.Bank.Transfer(ctx, req)
if err != nil {
return "", err
}
return resp.Reference, nil
}
ยังมีช่องว่างอยู่
ระหว่างขั้นที่ 1 กับ 2 มี race window อยู่ ถ้าเกิด retry พร้อมกัน 2 ครั้งพอดี
ก็ยังอาจโอนซ้ำได้ ในทางปฏิบัติควรเสริมด้วย lock ระดับ transaction
ในฐานข้อมูลของคุณเอง (เช่น insert แถวสถานะด้วย UNIQUE constraint ก่อนเรียกธนาคาร)
รูปแบบที่ 3 — Idempotent ที่ระดับฐานข้อมูล
สำหรับการเขียน ledger ของเราเอง วิธีที่แข็งแรงที่สุดคือให้ฐานข้อมูลบังคับ
-- ledger_entries มี UNIQUE (transaction_id, entry_type)
INSERT INTO ledger_entries (transaction_id, entry_type, wallet_id, amount_satang)
VALUES ($1, 'CREDIT', $2, $3)
ON CONFLICT (transaction_id, entry_type) DO NOTHING;
func (a *Activities) CreditWallet(
ctx context.Context, req TopUpRequest, bankRef string,
) (int64, error) {
tx, err := a.DB.BeginTx(ctx, nil)
if err != nil {
return 0, err
}
defer tx.Rollback()
_, err = tx.ExecContext(ctx, insertLedgerEntrySQL,
req.TransactionID, req.WalletID, req.AmountSatang)
if err != nil {
return 0, err
}
// อ่าน balance ที่คำนวณจาก ledger เสมอ
// เพื่อให้เรียกซ้ำแล้วได้ค่าเดิม
var balance int64
err = tx.QueryRowContext(ctx, selectBalanceSQL, req.WalletID).Scan(&balance)
if err != nil {
return 0, err
}
return balance, tx.Commit()
}
สังเกตว่า activity นี้ถูกเรียกซ้ำกี่ครั้งก็ได้ ผลลัพธ์เหมือนเดิมเสมอ — นี่คือนิยามของ idempotent ที่แท้จริง
Heartbeat
Heartbeat มี 3 หน้าที่ และทุกข้อสำคัญกับงานยาว:
- รับสัญญาณ cancel — cancellation ถูกส่งผ่าน heartbeat เท่านั้น activity ที่ไม่ heartbeat จะไม่มีวันรู้ว่าถูกยกเลิก
- ทำงานต่อจากจุดเดิม — ข้อมูลใน heartbeat อยู่รอดข้าม retry
- ตรวจจับ activity ที่แขวน — ถ้าหยุด heartbeat Temporal จะ timeout แล้ว retry
// ตัวอย่าง: จ่ายเงินเดือนพนักงาน 50,000 รายการ
func (a *Activities) BulkPayout(ctx context.Context, batchID string) (int, error) {
items, err := a.loadPayoutItems(ctx, batchID)
if err != nil {
return 0, err
}
// กู้ตำแหน่งจากความพยายามครั้งก่อน
start := 0
if activity.HasHeartbeatDetails(ctx) {
var last int
if err := activity.GetHeartbeatDetails(ctx, &last); err == nil {
start = last + 1
activity.GetLogger(ctx).Info("resuming batch", "from", start)
}
}
for i := start; i < len(items); i++ {
// ตรวจว่าถูกสั่งยกเลิกหรือยัง
if ctx.Err() != nil {
return i, ctx.Err()
}
if err := a.payoutOne(ctx, items[i]); err != nil {
return i, err
}
// บันทึกความคืบหน้า
activity.RecordHeartbeat(ctx, i)
}
return len(items), nil
}
ctx = workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
StartToCloseTimeout: 2 * time.Hour,
HeartbeatTimeout: 2 * time.Minute, // ต้องสั้นกว่า StartToClose มากๆ
})
ตั้ง HeartbeatTimeout ให้พอดี
สั้นเกินไป (เช่น 10 วินาที ทั้งที่แต่ละรายการใช้เวลา 15 วินาที) → activity ถูกฆ่าทิ้งทั้งที่ทำงานปกติ ยาวเกินไป → กว่าจะรู้ว่าแขวนก็เสียเวลาไปมาก ตั้งประมาณ 2–3 เท่าของช่วงเวลาที่คาดว่าจะ heartbeat จริง และจำไว้ว่าแต่ละ heartbeat นับเป็น action ที่มีค่าใช้จ่ายบน Temporal Cloud
Activity Context
func (a *Activities) ChargeBank(ctx context.Context, req TopUpRequest) (string, error) {
info := activity.GetInfo(ctx)
_ = info.Attempt // ครั้งที่เท่าไร (เริ่มที่ 1)
_ = info.WorkflowExecution.ID // workflow ID ที่เรียกมา
_ = info.ActivityID
_ = info.TaskQueue
_ = info.Deadline // deadline ของความพยายามครั้งนี้
logger := activity.GetLogger(ctx) // มี context ครบให้อัตโนมัติ
logger.Info("charging bank", "txn", req.TransactionID)
return "", nil
}
info.Attempt มีประโยชน์สำหรับ log และ metric แต่
ห้ามเอาไปประกอบเป็น idempotency key เพราะมันเปลี่ยนทุกครั้งที่ retry
Checklist ก่อนปล่อย activity ที่แตะเงิน
- ตั้ง
StartToCloseTimeoutยาวกว่า timeout ของ API ปลายทาง - ส่ง idempotency key ที่ deterministic ไปกับทุก request
- แยก error เชิงกฎธุรกิจให้เป็น non-retryable
- ตอบซ้ำได้ — เรียก 2 ครั้งได้ผลเหมือนเรียกครั้งเดียว
- ถ้ารันนานกว่า ~30 วินาที ต้องมี heartbeat และเช็ค
ctx.Err() - ไม่ส่งข้อมูลก้อนใหญ่เข้า/ออก (ส่ง reference แทน)
- ไม่ log เลขบัตร, PAN, OTP หรือข้อมูลระบุตัวตนลงใน log ธรรมดา
- มี unit test ที่ครอบคลุมเคสเรียกซ้ำ