跳到主要内容

Stream

Stream 是一种实时、可续读的消息传递通道,提供 best-effort storage。Step 或 Client 会立即写入消息,Client 使用 resume token 按顺序继续读取。

Stream 适合传递能够改善用户体验的 nice-to-have 数据,而不是关键业务数据:

  • LLM reasoning summary、生成文本 chunk 和 agent progress
  • 进度条、诊断事件与实时预览
  • Step 仍在运行时展示的临时状态

Stream 只提供 best-effort durability,底层使用 Redis 等内存存储。达到 capacity threshold 后触发 trim,或 storage backend 重启时,消息都可能丢失。需要 guaranteed durability 时,应使用 ChannelAttribute 等其他 primitive。

定义并注册 Stream

Stream 包含稳定名称、消息类型和 streamCapacityBytes。streamCapacityBytes 是相同 Flow type 和 Stream name 下所有 Flow execution 共享的近似容量,不是每个 Flow ID 独立的 limit。

下面的 Step 使用 buffered text writer。Dex 会在 Step result 之前 flush 累积的 progress。Handler 不会等待 Stream Store acknowledgement。

class RenderPreview(Step[str]):
def __init__(self, progress: Stream[str]) -> None:
self.progress = progress

async def execute(self, context: AsyncContext, input: str) -> StepDecision:
progress = self.progress.buffered_text(context)
progress.write(f"Rendering preview for {input}")
progress.write(f"Preview ready for {input}")
return graceful_complete(f"Rendered {input}")


class StreamFlow(Flow[str]):
progress = Stream("Progress", str, 10 * 1024 * 1024)

def get_persistence_schema(self) -> PersistenceSchema:
return PersistenceSchema.of(self.progress)

例子: examples/python/dex_examples/primitives/stream/stream_flow.py

语义

写入语义

Step 直接写入 Stream 时:

  • Dex 在 WaitForExecute 期间立即发送,不会缓存到 Step 完成。
  • Step 随后失败不会回滚已经写入的消息。
  • 一次 method invocation 可以向同一个或不同 Stream 写入任意数量的消息。每次写入都是 implicit heartbeat,并保留最近一次 explicit heartbeat value。
  • SDK 只确认本地编码并把消息交给 Worker output stream。Stream Store 被禁用、不可用、已满或拒绝写入时,消息会被丢弃,但 Step 不会失败。
  • 每条 Step message 的 source 都是 #StepExecutionID。同一个 Step execution 的不同 attempt 与 message 共用 source。Source 只表示来源,不要求唯一,也不会 deduplicate retry。

即使 Flow instance 不存在或已经结束,Client 仍可以写入。每次写入都要提供非空 source。Source 可以重复,也可以包含 #;每次调用都会追加一条消息。RPC handler 不能通过自己的 Context 写入,但可以通过注入的 Client dependency 发送 Stream message。

Python coroutine Step 调用 Stream.write 时不使用 await。Python synchronous Step 必须从 generator yield 返回的 StepOutput。其他 SDK 通过 Step Context 发送 frame,并在本地 enqueue 后返回。

Buffer text chunk

当许多小 chunk 共同组成一段文本时,例如 LLM delta,请使用 buffered text writer。默认 flush interval 是一秒,soft threshold 是 16 KiB UTF-8 数据。Timer、size threshold 或 invocation 完成时,都会发送当前非空 batch。Helper 会原样保留文本,不拆分单个 chunk,并忽略空 chunk。

Go、Java、Python async、TypeScript 和 Rust 会在发送最终 result 或 error 前停止 timer 并 flush 尾部。Python sync generator 使用 cooperative elapsed-time check,必须从显式的最终 flush yield output。空 buffer 不会产生 Stream message 或 implicit heartbeat。Retry 不会恢复尚未发送的文本;retry 前已经发送的 batch 可能再次出现。

可续读

ReadStream 每次返回一条消息,包括 value、resume token、created time 与 source。

  • 空 token 从当前保留的头部开始。
  • token 比当前头部更旧时,也从当前头部开始。
  • 下一次读取应原样传回上次返回的 token。
  • 没有下一条消息时,请求会 long-poll,直到消息到达或等待超时。

续读也是 best effort。读取太慢时,Client 可能错过已经被 trim 的消息。

列出消息

使用 ListStreamMessages 可以从最新到最早查看当前保留的消息。调用不会阻塞,而是立即返回不超过指定 page size 的消息。

before-page token 为空时,从当前保留的尾部开始。读取下一页时,应原样传回 next-page token;该 token 是排他锚点,并绑定到 Flow type、Flow ID 和 Stream name。next-page token 为空表示已经到达末页。

        page = await app_state.client.list_stream_messages(
required_query("workflowId"),
app_state.stream.progress,
required_int_query("pageSize"),
optional_query("beforePageToken", ""),
)

例子: examples/python/dex_examples/primitives/stream/controller.py

倒序分页提供 best-effort retained-message snapshot,并不是 transactional snapshot。第一页之后新追加的消息不会混入这个更早消息的分页链,但 trim 可能在两次调用之间删除消息。如果 trim 后 retained head 已经越过 page anchor,Dex 会返回空页,而不会从 tail 重新开始。

Server 要求 page size 为正数,并通过 maxReadMessages 设置上限。默认上限为 1000。

需要正向逐条消费、long-poll 和续读时,请使用 ReadStream。需要倒序、非阻塞地分页查看当前保留消息时,请使用 ListStreamMessages

近似容量

streamCapacityBytes 是相同 Flow type 和 Stream name 下所有 Flow execution 共享的近似容量,不是每个 Flow ID 独立的 limit。

Dex server 根据每条消息的 serialized value、Flow ID、source 和可配置的 per-message overhead 计算估算大小。

当 usage 达到 trim trigger threshold 时,全局只有一个 background trimmer 会删除所有 instance 中最旧的消息,直到 usage 降至 trim target threshold。极少数情况下,如果 trimming 太慢,usage 达到 capacity limit(streamCapacityBytes),write request 会被拒绝。

Server 还会限制单条 Stream message 的大小,默认上限为 100 KiB。

Best-effort storage

多 server 部署可以使用 Redis,单机本地 server 可以使用 memory backend。memory backend 会在进程重启时丢失全部内容。Redis 可以跨应用重启保留数据,但 trim、Redis 数据丢失、禁用 Stream Store 或容量压力仍可能删除或拒绝消息。