temporal

package module
v0.0.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 4, 2026 License: MIT Imports: 15 Imported by: 0

README

workflow/temporal

workflow/temporal 封装 Temporal Go SDK,提供客户端、Worker、默认消息工作流和可选 OpenTelemetry 追踪。模块保留 Temporal 的 Workflow、Activity、Signal、Query 等原生概念,适合长周期、可恢复、可观测的业务编排。

模块能力

  • 创建 Temporal 客户端,支持 HostPort 和 Namespace 配置。
  • 异步执行工作流,或同步等待结果。
  • 发送 Signal、执行 Query、取消工作流、查询执行描述。
  • 创建 Worker,并注册 Workflow 和 Activity。
  • 提供 StartSimpleWorker,快速处理 []byte 消息体。
  • 提供默认 BrokerMessageWorkflow,把消息体委托给 ProcessMessage Activity。
  • 通过 WithTracing() 为生产和消费路径创建 OpenTelemetry span。

安装

go get github.com/liujitcn/kratos-kit/workflow/temporal@latest

客户端

client, err := temporal.NewClient(
	temporal.WithClientHostPort("localhost:7233"),
	temporal.WithClientNamespace("default"),
)
if err != nil {
	return err
}
defer func() { _ = client.Close() }()

未传配置时默认连接 localhost:7233default namespace。

执行工作流

异步执行
runID, err := client.Execute(ctx, []byte(`{"order_id":"A1001"}`), temporal.ExecuteOptions{
	TaskQueue:  "order-task-queue",
	WorkflowID: "order-A1001",
})
if err != nil {
	return err
}

WorkflowFn 为空时使用内置 BrokerMessageWorkflow。该默认工作流会调用名为 ProcessMessage 的 Activity。

同步执行
result, err := client.ExecuteSync(ctx, []byte(`{"order_id":"A1001"}`), temporal.ExecuteOptions{
	TaskQueue:  "order-task-queue",
	WorkflowID: "order-A1001-sync",
})

ExecuteSync 等待工作流完成,并把结果读取为 []byte

自定义工作流
runID, err := client.Execute(ctx, order, temporal.ExecuteOptions{
	TaskQueue:  "order-task-queue",
	WorkflowID: "order-A1001",
	WorkflowFn: OrderWorkflow,
})

自定义 Workflow 和 Activity 需要在 Worker 上注册。

Worker

简单 Worker
worker, err := client.StartSimpleWorker(ctx, "order-task-queue",
	func(ctx context.Context, body []byte) error {
		return handleOrder(ctx, body)
	},
)
if err != nil {
	return err
}

StartSimpleWorker 会创建 Worker、注册默认消息处理 Activity 并立即启动。传入的 ctx 取消后会自动停止 Worker。

完整 Worker
worker, err := client.NewWorker(temporal.WorkerOptions{
	TaskQueue:  "order-task-queue",
	Workflows:  []any{OrderWorkflow},
	Activities: []any{ProcessOrder},
})
if err != nil {
	return err
}
if err := worker.Start(); err != nil {
	return err
}

NewWorker 会自动注册内置 BrokerMessageWorkflow,并额外注册 WorkerOptions.WorkflowsWorkerOptions.Activities

工作流控制

方法 说明
Signal(ctx, workflowID, runID, signalName, arg) 向运行中的工作流发送 Signal
Query(ctx, workflowID, runID, queryType, arg) 查询工作流状态
Cancel(ctx, workflowID, runID) 请求取消工作流
Describe(ctx, workflowID, runID) 获取工作流执行描述
TemporalClient() 返回底层 Temporal SDK Client

配置

ExecuteOptions
字段 说明
TaskQueue 工作流使用的任务队列
WorkflowID 工作流执行唯一标识
WorkflowFn 工作流函数,空值使用 BrokerMessageWorkflow
RunTimeout 单次运行超时
ExecutionTimeout 总执行超时,包含重试和 continue-as-new
TaskTimeout 单个 Workflow Task 超时
RetryPolicy Temporal 重试策略
CronSchedule Cron 调度表达式
IDReusePolicy WorkflowID 已存在时的复用策略
WorkerOptions
字段 说明
TaskQueue Worker 监听的任务队列
Options 原生 Temporal Worker 配置
Workflows 额外注册的 Workflow 函数列表
Activities 额外注册的 Activity 函数或结构体列表

OpenTelemetry

调用 client.WithTracing() 后,模块会使用 go.opentelemetry.io/otel 创建 tracer:

路径 Span 名称 Span Kind
Execute / ExecuteSync temporal-producer Producer
ProcessMessage Activity temporal-consumer Consumer

span 属性会包含 messaging.system=temporal 和当前 task queue。

本地开发

temporal server start-dev

默认 gRPC 地址为 localhost:7233,Web UI 地址为 http://localhost:8233

参考

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func BrokerMessageWorkflow

func BrokerMessageWorkflow(ctx workflow.Context, body []byte) error

BrokerMessageWorkflow 是默认消息工作流。 它接收消息体,并委托已注册的 ProcessMessage Activity 处理。 复杂编排场景可通过 WorkerOptions.Workflows 注册自定义工作流。

func WithClientHostPort

func WithClientHostPort(hostPort string) func(*ClientOptions)

WithClientHostPort 设置 Temporal Server 地址。

func WithClientNamespace

func WithClientNamespace(namespace string) func(*ClientOptions)

WithClientNamespace 设置 Temporal 命名空间。

Types

type ClientOptions

type ClientOptions struct {
	// HostPort 是 Temporal Server 地址,默认值为 "localhost:7233"。
	HostPort string

	// Namespace 是 Temporal 命名空间,默认值为 "default"。
	Namespace string
}

type ExecuteOptions

type ExecuteOptions struct {
	// TaskQueue 是工作流使用的任务队列。
	TaskQueue string

	// WorkflowID 是工作流执行的唯一标识。
	WorkflowID string

	// WorkflowFn 是需要执行的工作流函数,空值时使用默认工作流。
	WorkflowFn any

	// RunTimeout 是单次工作流运行的最大耗时。
	RunTimeout time.Duration

	// ExecutionTimeout 是包含重试和 continue-as-new 在内的总执行超时。
	ExecutionTimeout time.Duration

	// TaskTimeout 是单个工作流任务超时。
	TaskTimeout time.Duration

	// RetryPolicy 是工作流重试策略。
	RetryPolicy *temporal.RetryPolicy

	// CronSchedule 是工作流定时调度表达式。
	CronSchedule string

	// IDReusePolicy 控制 WorkflowID 已存在时的复用行为。
	IDReusePolicy enums.WorkflowIdReusePolicy
}

type WorkerOptions

type WorkerOptions struct {
	// TaskQueue 是 Worker 监听的任务队列。
	TaskQueue string

	// Options 是原生 Temporal Worker 配置。
	Options worker.Options

	// Workflows 是需要额外注册的工作流函数列表。
	Workflows []any

	// Activities 是需要额外注册的 Activity 函数或结构体列表。
	Activities []any
}

type WorkflowClient

type WorkflowClient struct {
	// contains filtered or unexported fields
}

WorkflowClient 提供 Temporal 工作流操作的高层封装。

func NewClient

func NewClient(opts ...func(*ClientOptions)) (*WorkflowClient, error)

NewClient 创建并连接 Temporal 客户端。

func (*WorkflowClient) Cancel

func (wc *WorkflowClient) Cancel(ctx context.Context, workflowID, runID string) error

Cancel 请求取消正在运行的工作流。

func (*WorkflowClient) Close

func (wc *WorkflowClient) Close() error

Close 关闭底层 Temporal 客户端连接。

func (*WorkflowClient) Describe

Describe 获取工作流执行描述信息。

func (*WorkflowClient) Execute

func (wc *WorkflowClient) Execute(ctx context.Context, args any, opts ExecuteOptions) (string, error)

Execute 异步启动一次工作流执行。 返回本次执行的 run ID。

func (*WorkflowClient) ExecuteSync

func (wc *WorkflowClient) ExecuteSync(ctx context.Context, args any, opts ExecuteOptions) ([]byte, error)

ExecuteSync 启动工作流并等待完成,返回工作流结果。

func (*WorkflowClient) NewWorker

func (wc *WorkflowClient) NewWorker(opts WorkerOptions) (*WorkflowWorker, error)

NewWorker 创建 Temporal Worker。 创建后不会自动启动,需要调用 Start 开始轮询任务。

func (*WorkflowClient) Query

func (wc *WorkflowClient) Query(ctx context.Context, workflowID, runID, queryType string, arg any) (any, error)

Query 查询正在运行的工作流状态。

func (*WorkflowClient) Signal

func (wc *WorkflowClient) Signal(ctx context.Context, workflowID, runID, signalName string, arg any) error

Signal 向正在运行的工作流发送信号。

func (*WorkflowClient) StartSimpleWorker

func (wc *WorkflowClient) StartSimpleWorker(ctx context.Context, taskQueue string, handler func(ctx context.Context, body []byte) error, opts ...func(*WorkerOptions)) (*WorkflowWorker, error)

StartSimpleWorker 创建带单一消息处理函数的 Worker。 这是从任务队列开始处理消息的最简方式。

func (*WorkflowClient) TemporalClient

func (wc *WorkflowClient) TemporalClient() client.Client

TemporalClient 返回底层 Temporal SDK 客户端,供高级场景使用。

func (*WorkflowClient) WithTracing

func (wc *WorkflowClient) WithTracing()

WithTracing 为客户端启用 OpenTelemetry 链路追踪。

type WorkflowWorker

type WorkflowWorker struct {
	sync.RWMutex
	// contains filtered or unexported fields
}

WorkflowWorker 管理负责轮询任务的 Temporal Worker。

func (*WorkflowWorker) IsRunning

func (ww *WorkflowWorker) IsRunning() bool

IsRunning 返回 Worker 是否已经成功启动且尚未停止。

func (*WorkflowWorker) RegisterActivity

func (ww *WorkflowWorker) RegisterActivity(fn any)

RegisterActivity 注册 Activity 函数或结构体。 必须在 Start 之前调用。

func (*WorkflowWorker) RegisterWorkflow

func (ww *WorkflowWorker) RegisterWorkflow(fn any)

RegisterWorkflow 注册工作流函数。 必须在 Start 之前调用。

func (*WorkflowWorker) Start

func (ww *WorkflowWorker) Start() error

Start 启动 Worker 轮询任务。

func (*WorkflowWorker) Stop

func (ww *WorkflowWorker) Stop()

Stop 优雅停止 Worker。

func (*WorkflowWorker) TaskQueue

func (ww *WorkflowWorker) TaskQueue() string

TaskQueue 返回当前 Worker 监听的任务队列名。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL