高级:多父 Flow 分区
将请求路由到多个父 Flow ID,使活跃 SubFlow 总数能够超过单个父 Flow 的并发限制。
使用稳定的分区键,例如请求 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