跳到主要内容

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

字段作用
TimeoutFlow timeout。省略或设为零可禁用。
TimeoutPolicyFlow timeout 触发时的行为。
TimeoutHandlerOptionsTimeout handler 的 execution、retry、recovery、durability、locking 与 state loading。
IDReusePolicy启动时如何复用已有 Flow ID。
RequestID重试同一次启动请求的幂等键。
AlreadyStartedOptions如何处理已启动的 Flow。
RetryPolicy单次 Flow execution 失败后的重试策略。
StartDelaystart 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 表示这次调用是否在重试同一次启动请求。

这些选项如何一起工作

这些选项依次回答不同的问题:

  1. RequestID 标识一次 start request。省略时,SDK 会生成新的 UUID,所以另一次 StartFlow 调用是新的 request。只在重试完全相同的 start 时复用它。Dex 比较的是 RequestID,不比较 start input 或 options。
  2. IDReusePolicy 决定这个 request 能否用该 Flow ID 启动另一个 run。
  3. AlreadyStartedOptions.IgnoreError 只在 ID reuse policy 因 already-started 拒绝启动时生效。只有 RequestID 相同时,它才返回已有的 Run ID。
重试请求AlreadyStartedOptions.IgnoreErrorID 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。InitialIntervalBackoffCoefficientMaximumInterval 设置指数退避。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。
AttributeStoreNamesDex Server Attribute Store 名称,供 opt-in attribute sync 的 Attribute 使用。

用 Client 更新 FlowConfig

Flow active 时可以调用 UpdateFlowConfig。只发送要修改的字段。它们会与已有配置合并。