
Temporal Go 实施手册本手册为使用 Temporal Go SDK 实现持久编排提供了生产就绪的模式和深入的技术指导。目录确定性戒律工作流版本管理活动设计与幂等性面向规模的工作器配置上下文与心跳拦截器与可观测性1. 确定性戒律在 Go 中工作流是必须完全相同地重放的状态机。违反这些规则会导致确定性不匹配错误。❌ 绝不使用原生 Go 并发错误:go myFunc()正确:workflow.Go(ctx, func(ctx workflow.Context) { ... })原因:workflow.Go允许 Temporal 编排器在重放期间跟踪和暂停 goroutine。❌ 绝不使用原生时间错误:time.Now()、time.Sleep(d)、time.After(d)正确:workflow.Now(ctx)、workflow.Sleep(ctx, d)、workflow.NewTimer(ctx, d)❌ 绝不确定地迭代 map错误:for k, v : range myMap { ... }正确:收集键、排序然后迭代。keys:make([]string,0,len(myMap))fork:rangemyMap{keysappend(keys,k)}sort.Strings(keys)for_,k:rangekeys{v:myMap[k];...}❌ 绝不直接执行外部 I/O错误:在工作流内部执行http.Get(https://api.example.com)或os.ReadFile(data.txt)。正确:将所有 I/O 封装在活动中并使用workflow.ExecuteActivity调用它。原因:外部调用是不确定的其结果会在重放之间变化。❌ 绝不使用不确定的随机数错误:在工作流内部使用rand.Int()、uuid.New()。正确:将随机种子或 UUID 作为工作流输入参数传入或在活动中生成它们。原因:rand.Int()在每次重放时产生不同的值导致确定性不匹配。2. 工作流版本管理当你需要在运行中的工作流中更改逻辑时必须使用workflow.GetVersion。模式安全的逻辑更新constVersionV21funcMyWorkflow(ctx workflow.Context)error{v:workflow.GetVersion(ctx,ChangePaymentStep,workflow.DefaultVersion,VersionV2)ifvworkflow.DefaultVersion{// Old logic: kept alive until all pre-existing workflow runs complete.returnworkflow.ExecuteActivity(ctx,OldActivity).Get(ctx,nil)}// New logic: all new and resumed workflow runs use this path.returnworkflow.ExecuteActivity(ctx,NewActivity).Get(ctx,nil)}模式完全迁移后的清理一旦确认没有正在运行的工作流实例处于DefaultVersion通过 Temporal Web UI 或tctl验证就可以安全地移除旧分支funcMyWorkflow(ctx workflow.Context)error{// Pin minimum version to V2; histories from before the migration will// fail the determinism check (replay error) if they replay against this code.// Only remove the old branch after confirming zero running instances on DefaultVersion.workflow.GetVersion(ctx,ChangePaymentStep,VersionV2,VersionV2)returnworkflow.ExecuteActivity(ctx,NewActivity).Get(ctx,nil)}3. 活动设计与幂等性活动可能执行多次。它们必须是幂等的。模式使用 Upsert 而非 Insert不要使用简单的INSERT而是使用UPSERT或带有幂等键如WorkflowID或RunID的先检查后执行。func(a*Activities)ProcessPayment(ctx context.Context,req PaymentRequest)error{info:activity.GetInfo(ctx)// Use info.WorkflowExecution.ID as part of your idempotency key in DBreturna.db.UpsertPayment(req,info.WorkflowExecution.ID)}4. 面向规模的工作器配置优化的工作器选项w:worker.New(c,task-queue,worker.Options{MaxConcurrentActivityExecutionSize:100,// Limit based on resource constraintsMaxConcurrentWorkflowTaskExecutionSize:50,WorkerActivitiesPerSecond:200,// Rate limit for this worker clusterWorkerStopTimeout:time.Minute,// Allow activities to finish})5. 上下文与心跳传播元数据使用工作流拦截器或自定义Header传播沿调用链传递追踪 ID 或用户身份。活动心跳对于长时间运行的活动为了在StartToCloseTimeout到期前检测工作器崩溃心跳是强制性的。funcLongRunningActivity(ctx context.Context)error{fori:0;i100;i{activity.RecordHeartbeat(ctx,i)// Report progressselect{case-ctx.Done():returnctx.Err()// Handle cancellationdefault:// Do work}}returnnil}6. 拦截器与可观测性自定义工作流拦截器使用拦截器注入结构化日志Zap/Slog或执行全局错误分类。拦截器必须通过 Temporal 为每个工作流任务实例化的根WorkerInterceptor进行装配。// Step 1: Implement the root WorkerInterceptor (registered on worker.Options)typeMyWorkerInterceptorstruct{interceptor.WorkerInterceptorBase}func(w*MyWorkerInterceptor)InterceptWorkflow(ctx workflow.Context,next interceptor.WorkflowInboundInterceptor,)interceptor.WorkflowInboundInterceptor{returnmyWorkflowInboundInterceptor{next:next}}// Step 2: Implement the per-workflow inbound interceptortypemyWorkflowInboundInterceptorstruct{interceptor.WorkflowInboundInterceptorBase next interceptor.WorkflowInboundInterceptor}func(i*myWorkflowInboundInterceptor)ExecuteWorkflow(ctx workflow.Context,input*interceptor.ExecuteWorkflowInput,)(interface{},error){workflow.GetLogger(ctx).Info(Workflow started,type,workflow.GetInfo(ctx).WorkflowType.Name)result,err:i.next.ExecuteWorkflow(ctx,input)iferr!nil{workflow.GetLogger(ctx).Error(Workflow failed,error,err)}returnresult,err}// Step 3: Register on the workerw:worker.New(c,task-queue,worker.Options{Interceptors:[]interceptor.WorkerInterceptor{MyWorkerInterceptor{}},})应避免的反模式巨型工作流:在单个工作流中保留过多状态。如果事件历史超过 5 万条事件请使用ContinueAsNew。臃肿的活动:在活动内部执行编排。活动应该是工作单元。全局变量:在工作流中使用全局变量。它们不会在工作器重启后保留。工作流中的原生并发:使用go例程、mutexes或channels会在重放期间导致竞态条件和确定性错误。7. SideEffect 和 MutableSideEffect当你需要一个仅捕获一次并完全相同重放的非确定性值时——例如在工作流内部生成 UUID 或读取一次性配置快照——使用workflow.SideEffect。// SideEffect: called only on first execution; result is recorded in history and// replayed deterministically on all subsequent replays.// Requires: go.temporal.io/sdk/workflowencodedID:workflow.SideEffect(ctx,func(ctx workflow.Context)interface{}{returnuuid.NewString()})varrequestIDstringiferr:encodedID.Get(requestID);err!nil{returnerr}何时使用MutableSideEffect: 当值可能在工作流任务之间变化但每个历史事件仍必须是确定性的例如在工作流运行期间更新的功能标志时。// MutableSideEffect: re-evaluated on each workflow task, but only recorded in// history when the value changes from the previous recorded value.encodedFlag:workflow.MutableSideEffect(ctx,feature-flag-v2,func(ctx workflow.Context)interface{}{returnfeatureFlagEnabled// read from workflow-local state, NOT an external call},func(a,binterface{})bool{returna.(bool)b.(bool)},)varenabledboolencodedFlag.Get(enabled)警告:不要使用SideEffect作为在工作流内部调用外部 APIHTTP、数据库的变通方法。所有外部 I/O 仍必须通过活动进行。