等待并流式接收更新
当调用方需要接收一系列更新,而不是一个就绪边界时,使用 Stream。ReadStream 会等待所提供 resume token 之后的下一条消息。响应带有一个新 token;在下一次读取中传入它,即可按顺序继续。
StreamFlow 示例公开一个 progress Stream。它的 read endpoint 最多长轮询 20 秒,返回一条消息,并提供客户端应保存的 token。客户端处理这条消息后可以立刻发起下一次读取。
Stream 是 best-effort、有序消息,而不是最终 Flow 结果。当调用方必须在继续前等待一个特定 durable 条件时,使用 WaitForStepCompletion 或 WaitForAttributeMatch。
核心实现
@blueprint.get("/read")
async def read() -> Response:
message = await app_state.client.read_stream(
required_query("workflowId"),
app_state.stream.progress,
optional_query("resumeToken", ""),
timedelta(seconds=20),
)
return jsonify(
value=message.value,
resume_token=message.resume_token,
created_time=message.created_time.isoformat(),
source=message.source,
)
例子: examples/python/dex_examples/primitives/stream/controller.py