跳到主要内容

高级:多父 Flow 分区

将请求路由到多个父 Flow ID,使活跃 SubFlow 总数能够超过单个父 Flow 的并发限制。

Definition graph

SubmitRequestFlow

Valid

python · examples/python/dex_examples/patterns/parallel-subflows/submit_request_flow.py

使用稳定的分区键,例如请求 ID 哈希对父 Flow 数量取模。相同键的请求始终到达同一个父 Flow。每个父 Flow 独立执行 Channel 缓冲区和 SubFlow 并发限制,因此增加父 Flow 数量即可提高总容量,而不会让单个 Flow execution 无界增长。

改变分区数量或顺序会重新映射键。扩容时如果需要稳定亲和性,请使用带版本的父 Flow 集合或一致性哈希。

核心实现

提交方先调用目标父 Flow 的 RPC。如果父 Flow 不存在或已不活跃,就以该请求作为初始输入启动父 Flow。如果另一个提交方抢先完成启动,本次启动会失败,并由 SubmitStep 重试。下一次执行会在已活跃的父 Flow 上调用 RPC。初始请求不会丢失,因为父 Flow 的 InitStep 会把它发布到 RequestChannel

async def enqueue_request(client, parent_flow, parent_id: str, request: str) -> bool:
try:
return await client.invoke_rpc(parent_flow.send_request, parent_id, request)
except FlowNotActiveError:
await client.start_flow(
parent_flow,
parent_id,
ParentInput([request], DEFAULT_CONCURRENCY),
StartFlowOptions(id_reuse_policy=IdReusePolicy.ALLOW_IF_NOT_RUNNING),
)
return True


def partition(request: str, partitions: int) -> int:
hash_value = 2_166_136_261
for byte in request.encode():
hash_value ^= byte
hash_value = hash_value * 16_777_619 & 0xFFFFFFFF
return hash_value % partitions

例子: examples/python/dex_examples/patterns/parallel-subflows/submit_request_flow.py

相关内容