Flow options
Dex 允许你自定义 Flow execution 的行为。
- StartFlowOptions 在客户端启动 Flow(或 SubFlow)时生效,启动后不可变。
- 它们控制 timeout、ID reuse、初始 Attribute、request 幂等,以及该次 execution 的可选 config override。
- FlowConfig 提供运行时行为,启动后可变。
- 包括 Step durability 默认值、Worker 路由、active-Step 搜索索引、continue-as-new 阈值与 Attribute Store projection target。
StartFlowOptions
| 字段 | 作用 |
|---|---|
| Timeout | Flow timeout。省略或设为零可禁用。 |
| TimeoutPolicy | Flow timeout 触发时的行为。 |
| TimeoutHandlerOptions | Timeout handler 的 execution、retry、recovery、durability、locking 与 state loading。 |
| IDReusePolicy | 启动时如何复用已有 Flow ID。 |
| RequestID | 重试同一次启动请求的幂等键。 |
| AlreadyStartedOptions | 如何处理已启动的 Flow。 |
| RetryPolicy | 单次 Flow execution 失败后的重试策略。 |
| StartDelay | start Step 运行前的延迟。 |
| Attributes | 初始 Attribute 值。 |
| FlowConfig | 每次 execution 的 Flow 配置。 |
如需在 Flow 截止时间执行应用定义的处理,请参阅 Graceful Timeout。
def example_start_flow_options() -> StartFlowOptions:
return StartFlowOptions(
timeout=timedelta(minutes=30),
timeout_policy=FlowTimeoutPolicy.HANDLER,
start_delay=timedelta(minutes=5),
id_reuse_policy=IdReusePolicy.DISALLOW,
retry_policy=RetryPolicy(
initial_interval=timedelta(minutes=1),
backoff_coefficient=2,
maximum_interval=timedelta(minutes=10),
maximum_attempts=3,
),
config_override=FlowConfig(step_durability=StepDurability.SYNC),
ignore_already_started=True,
request_id="start-order-123",
).with_attribute(status, "queued")
例子: examples/python/dex_examples/primitives/flow/controller.py
Timeout 与 TimeoutPolicy
Timeout 是整个 Flow 的一个 durable timer。只有 Timeout 为正数时才能设置 TimeoutPolicy。
- FAIL — 以 timeout error 失败 Flow。RetryPolicy 可以重试它。
- CANCEL — 取消 Flow,不会重试。
- HANDLER — 运行一次逻辑 Flow timeout-handler execution。
提供正数 Timeout 时,Dex SDK 默认使用 FAIL。Flow 有 timeout handler 时,则使用 HANDLER。没有 handler 时使用 HANDLER 无效。一次 retry 会获得新的 timeout budget。
只有 Timeout 为正数且最终 policy 为 HANDLER 时,才能设置 TimeoutHandlerOptions。它采用 Execute 语义:method timeout 限制单次 attempt,Retry 可以运行多次 attempt,HeartbeatTimeout 检测停滞的 regular attempt,Failure 可以在 retry 耗尽后进入已注册的无输入 Step。Recovery Step 会收到 null 或 unit input,并从 handler Context 读取最终 error。
Handler 的 attempt 与 retry 计时从 soft Flow timeout 触发时开始,因此可以超过原始 Flow deadline。未配置 failure target 时,handler retry 耗尽会使 Flow 失败。
普通 Attribute 和所有 Channel size metadata 会自动加载。AttributeMap value 与 pending Channel message 必须通过 TimeoutHandlerOptions 显式选择。Handler 可以先暂存 Attribute write、Channel deletion 与 Channel publication,再返回 decision。Continue-as-new 和 Flow retry 都会保留这些 options;Flow retry 会开始新的 soft-timeout budget。
把 TimeoutPolicy 设为 HANDLER,timeout 触发时就会运行下面的代码。这个 handler 记录 Flow 结束的原因,再用应用自己的原因失败该 Flow。
def handle_timeout(self, context: Context) -> StepDecision:
status.set(context, "timed out")
return force_fail("processing deadline reached")
例子: examples/python/dex_examples/primitives/flow/example_flow.py
Flow ID、RequestID 与 AlreadyStartedOptions
IDReusePolicy 决定这个 Flow ID 已存在时该怎么做。RequestID 表示这次调用是否在重试同一次启动请求。
这些选项如何一起工作
这些选项依次回答不同的问题:
- RequestID 标识一次 start request。省略时,SDK 会生成新的 UUID,所以另一次 StartFlow 调用是新的 request。只在重试完全相同的 start 时复用它。Dex 比较的是 RequestID,不比较 start input 或 options。
- IDReusePolicy 决定这个 request 能否用该 Flow ID 启动另一个 run。
- AlreadyStartedOptions.IgnoreError 只在 ID reuse policy 因 already-started 拒绝启动时生效。只有 RequestID 相同时,它才返回已有的 Run ID。
| 重试请求 | AlreadyStartedOptions.IgnoreError | ID reuse policy 拒绝启动时 |
|---|---|---|
| 不提供 RequestID,或提供新的 RequestID | 关闭或开启 | SDK 发送新的 UUID。Dex 返回 already-started error。 |
| 相同 RequestID | 关闭 | Dex 返回 already-started error。 |
| 相同 RequestID | 开启 | Dex 返回已有 Run ID。 |
ID reuse policy 决定 Dex 是否会首先拒绝启动:
| IDReusePolicy | 已有 active Flow | 已有 closed Flow |
|---|---|---|
| IDReuseDefault | 拒绝启动。RequestID 相同且启用 IgnoreError 时返回 active Run ID。 | 启动新的 run。 |
| IDReuseAllowIfPreviousFailed | 拒绝启动。RequestID 相同且启用 IgnoreError 时返回 active Run ID。 | 只有上一次 run 失败后才启动新的 run;否则拒绝。 |
| IDReuseAllowIfNotRunning | 拒绝启动。RequestID 相同且启用 IgnoreError 时返回 active Run ID。 | 启动新的 run。 |
| IDReuseDisallow | 拒绝启动。RequestID 相同且启用 IgnoreError 时返回已有 Run ID。 | 拒绝启动。RequestID 相同且启用 IgnoreError 时返回已有 Run ID。 |
| IDReuseTerminateIfRunning | 终止 active run 并启动 replacement。不会产生 already-started error。 | 启动新的 run。 |
Dex Server 的逻辑
Dex Server 要求 RequestID 非空。省略时,SDK 会在发送请求前生成它。服务端将它记录在 Flow execution 中,再使用请求的 ID reuse policy 让已配置的 backend 启动 Flow。IDReuseDefault 使用服务端默认行为:没有 active execution 时启动新的 Flow execution。
如果 backend 返回 Flow 已启动,Dex 会检查 AlreadyStartedOptions.IgnoreError。启用时,Dex 读取已有 run 记录的 RequestID。只有两个 RequestID 相同时才返回该 Run ID;否则仍然返回 already-started error。
RetryPolicy
RetryPolicy 会在一次 Flow execution 失败后重试它。MaximumAttempts 包含首次 execution。InitialInterval、BackoffCoefficient 与 MaximumInterval 设置指数退避。FAIL timeout 可以使用这个 retry policy;CANCEL timeout 不可以。
StartDelay
StartDelay 会立即接受要启动的 Flow,但在 StartStep 可以运行前等待。它适合一次性延迟,例如新订单创建五分钟后再处理。周期性工作请使用 durable Timer 循环。
Attributes
Attributes 写入初始值。每个值必须匹配 Flow persistence schema 中的 Attribute 或 Attribute map。
FlowConfig
在 StartFlowOptions 中使用 ConfigOverride 修改一个 Flow 的初始配置。它只修改提供的字段,不会整体替换 FlowConfig。
启动时,Dex 会把提供的字段覆盖到 Dex Server 默认值上。省略的字段保留默认值。
| 字段 | 行为 |
|---|---|
| ActiveStepSearchMode | 索引所有 active Step、只索引运行过 WaitFor 的 Step,或不索引 active Step。默认只索引带 WaitFor 的 Step。 |
| ContinueAsNewThreshold | 追踪到这个事件数后 continue as new。零会关闭它。 |
| ContinueAsNewPageSizeBytes | 限制 continue-as-new 带入的每一页 state。零使用 1 MiB。 |
| StepDurability | 没有 Step override 的 method 使用的默认 durability。Dex 依次使用 method override、这里的值,最后使用 SYNC。 |
| WorkerTarget | 用于 Step 和 RPC invocation 的 WorkerService endpoint。 |
| AttributeStoreNames | Dex Server Attribute Store 名称,供 opt-in attribute sync 的 Attribute 使用。 |
用 Client 更新 FlowConfig
Flow active 时可以调用 UpdateFlowConfig。只发送要修改的字段。它们会与已有配置合并。