BLOGE
0.9.8-RC1在线版 · 事实校验 2026-09-15 · English
附录 C —— Streaming
普通 graph 会在上游节点完全结束后才把输出交给下游。Streaming 改变了这个契约:数据在上游仍在产出的过程中,通过 channel 逐条流向下游。
什么场景需要 Streaming
大多数 BLOGE graph 采用"批量完成"模式:每个节点执行完毕、输出被存储,之后下游节点才能看到该输出。当数据集有界且可以舒适地放进内存时,这种模式完全够用。
以下场景需要 streaming:
- 上游增量产出数据 —— 音频块、数据库游标分页、实时传感器读数——你希望一边收到一边处理。
- 数据集大到一次性物化很浪费 —— 你更希望通过有界缓冲区做流水线。
- 延迟比吞吐量更重要 —— 在剩余数据还在计算时就向用户展示部分结果。
如果以上都不适用,普通节点和 foreach 更简单,也更容易测试。如果你真正困难的不是数据量,而是生命周期 —— 多轮外部交互、命名状态或信号所有权 —— 那么应该继续阅读第 14 章到第 16 章,而不是优先选择 streaming 原语。
三种 Streaming 原语
BLOGE 提供三种产出流式输出的构造,各服务于不同模式。
1. stream node —— 流式源或变换
stream node 声明一个 operator,其输出通过内部 NodeChannel 逐条发射,而非作为单个值返回。
stream node audioCapture : AudioCapture {
buffer = 64
}
Operator 按自身节奏向 channel 推送条目。下游节点可以实时消费 channel,也可以等待 channel 关闭后接收物化后的完整列表。
适用场景: 单个 operator 产出无界或长时运行的条目序列(音频采集、事件流、分页 API)。
2. stream foreach —— 遍历集合并流式输出
stream foreach 是 foreach 的流式变体。它遍历集合,每个元素处理完成后立即将结果发射到下游,而非等所有元素都处理完。
stream foreach processOrders : item in loadOrders.output.orders {
buffer = 16
node processItem : OrderProcessor {
input {
order = item
}
}
}
适用场景: 集合是有界的,但你希望下游节点在最后一个元素处理完之前就开始消费结果。