บทที่ 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 WorkflowNexus
ขอบเขต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 และ cancelAsync Activity Completion
หลาย activity ต้องอยู่เครื่องเดียวกันSession (หรือรวมเป็น activity เดียว)
เรียกบริการของทีมอื่นNexus
งานเดี่ยวที่อยากได้ retry ของ TemporalStandalone Activity (ระวัง preview)

ข้อควรระวังร่วมกันของทั้งบท

pattern ในบทนี้แลกความทนทานบางส่วนกับความเร็วหรือความยืดหยุ่นเกือบทุกตัว ก่อนหยิบมาใช้ ให้ถามเสมอว่า "ถ้าส่วนนี้ล้มกลางคัน เงินของลูกค้าจะอยู่ตรงไหน" ถ้าตอบไม่ได้ชัดเจน ให้กลับไปใช้ activity ปกติกับ saga จากบทที่ 6