บทที่ 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 ตามลักษณะงาน

ActivityPolicy ที่แนะนำเหตุผล
เรียก bank API (หักเงิน/โอน)retry ไม่จำกัด + backoff + idempotency keyธนาคารล่มเป็นเรื่องชั่วคราว ยอมรอดีกว่าล้มธุรกรรม
ตรวจ limit / KYCretry ปกติ แต่ error เชิงกฎธุรกิจต้อง non-retryable"เกิน limit" retry อีกล้านครั้งก็ไม่ผ่าน
เขียน ledgerretry ไม่จำกัดDB ล่มคือปัญหาชั่วคราว และห้ามยอมให้ขั้นนี้ล้มเด็ดขาด
ส่ง push / SMSMaximumAttempts: 3ไม่กระทบเงิน ล้มแล้วปล่อยผ่านได้
เช็คสถานะจากธนาคาร (polling)BackoffCoefficient: 1 + InitialInterval = ช่วงที่อยากถามทำให้ retry กลายเป็น polling loop ที่ไม่กิน history — ดูบทที่ 13

แยก error ให้ถูกประเภท

นี่คือการตัดสินใจที่สำคัญที่สุดในการเขียน activity — ผิดแล้วเจ็บทั้ง 2ทาง

RetryableNon-retryable
network error, connection resetinput ไม่ถูกต้อง (validation)
timeout, 503, 502authentication / 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.IDtime.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 หน้าที่ และทุกข้อสำคัญกับงานยาว:

  1. รับสัญญาณ cancel — cancellation ถูกส่งผ่าน heartbeat เท่านั้น activity ที่ไม่ heartbeat จะไม่มีวันรู้ว่าถูกยกเลิก
  2. ทำงานต่อจากจุดเดิม — ข้อมูลใน heartbeat อยู่รอดข้าม retry
  3. ตรวจจับ 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 ที่ครอบคลุมเคสเรียกซ้ำ