Skip to main content

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 foreachforeach 的流式变体。它遍历集合,每个元素处理完成后立即将结果发射到下游,而非等所有元素都处理完。

stream foreach processOrders : item in loadOrders.output.orders {
buffer = 16
node processItem : OrderProcessor {
input {
order = item
}
}
}

适用场景: 集合是有界的,但你希望下游节点在最后一个元素处理完之前就开始消费结果。

3. stream loop —— 轮询并流式输出

stream looploop 的流式变体。每次迭代的输出会立即发射到下游,循环持续到 until 满足或达到 max_iterations

stream loop checkStatus {
max_iterations = 100
delay = 2s
depends_on = [initMonitor]
buffer = 8
node pollStatus : StatusPoller {
input {
serviceId = initMonitor.output.serviceId
iteration = loopIteration
}
}
until pollStatus.output.status == "ready"
}

适用场景: 你在轮询或持续产出条目,希望立即转发每个结果,而不是收集所有迭代后再一次性交付。


.stream.output —— 两种边类型

当下游节点引用流式节点的结果时,边的类型决定了数据何时以及如何到达。

语法边类型行为
upstream.streamStreamEdge下游实时接收条目,逐条到达。下游 operator 必须能处理增量输入。
upstream.outputDirectEdge下游等待流关闭后,接收完整的物化 List<T>。适用于下游 operator 需要全部数据才能开始工作的场景。

如何选择

问自己:下游 operator 是否需要所有数据才能开始?

  • → 使用 .output。例如:报告生成器需要所有已处理的订单来计算汇总。
  • → 使用 .stream。例如:语音转文字节点在音频块到达时即开始转录。

voice pipeline 示例在同一个 graph 中展示了两种边:

/// 通过 StreamEdge 实时接收音频块
stream node speechToText : SpeechToText {
input {
audio = audioCapture.stream // StreamEdge —— 实时转发
}
buffer = 16
}

/// 通过 DirectEdge 等待完整转录
node textAnalysis : TextAnalyzer {
depends_on = [speechToText]
input {
transcript = speechToText.output // DirectEdge —— 物化 List<T>
}
}

缓冲区大小

每个流式构造都有 buffer 参数,用于设置内部 NodeChannel 环形缓冲区的容量。

stream node audioCapture : AudioCapture {
buffer = 64 // 环形缓冲区最多持有 64 个条目
}

默认值: 16 个条目。

如何思考缓冲区大小

缓冲区解耦了生产者速率和消费者速率。当缓冲区满时,生产者会阻塞(背压),直到消费者追上来。

场景建议
生产者远快于消费者增大缓冲区以吸收突发流量;差距过大时背压仍然会启动。
生产者和消费者速率接近默认值 16 通常够用。
内存受限或条目较大减小缓冲区以限制内存使用;接受更频繁的背压暂停。
延迟敏感的流水线保持缓冲区较小,避免条目排队时间过长。

没有万能公式——从默认值开始,测量,再调整。当可观测性监听器激活时,引擎会记录背压事件,调优变得很直观。


Back-pressure、取消与持久化边界

事件契约要测试什么
Buffer 已满Producer 等待,不能静默丢 item放慢 consumer,观察内存有界且最终继续
下游取消取消向上游传播并关闭 channel断言 producer 停止,取消后不再发射
Producer 失败Channel 以失败关闭,物化 .output 不能伪装完整断言下游失败与部分输出策略
进程崩溃内存 channel 内容不是 durable checkpoint重启并证明哪些 item 来自 source 或 store 重放

Streaming 与 durable execution 解决不同问题。有界 channel 控制运行进程里的实时流动,不会让途中 item 自动跨崩溃保存。需要 replay 时,必须说出 durable source offset 或 checkpoint,并定义重复处理。不要从 buffer 大小推导 exactly-once delivery。


完整示例

仓库中以下文件用真实场景演示了三种 streaming 原语。

文件原语展示内容
streaming-batch.blogestream foreach流式处理客户订单;下游报告使用 .output(DirectEdge)物化全部结果。
streaming-status-monitor.blogestream loop每 2 秒轮询服务状态,将每次状态更新流式转发到下游;服务就绪时提前退出循环。
voice-pipeline.blogestream node音频采集 → 语音转文字 → 文本分析。在同一个 graph 中展示 .stream(实时)和 .output(物化)两种边。

打开这些文件,阅读文档注释,追踪数据如何从流式源经过缓冲区流向下游消费者。修改 buffer 值,观察背压行为的变化。


速查表

Diagram: appendix-c-streaming figure 1


常见错误

错误为什么会犯如何纠正
在 operator 需要全部数据时使用 .stream习惯性地选择"更快的路径"如果 operator 在看到所有条目之前无法产出输出,就使用 .output
把缓冲区设得很大"以防万一"想完全避免背压过大的缓冲区会延迟背压信号并浪费内存。从默认值开始。
忘了 stream foreach 同时也在收集 .output以为 streaming 意味着什么都不会物化条目被流式转发的同时也在收集;.output 会在流关闭后提供完整列表。
不测试背压路径只用小数据集测试,缓冲区从未被填满在测试中用 MockOperator.delaying(...) 模拟慢消费者,验证 graph 在背压下的行为。

与主线章节的关系

Streaming 没有被安排在单独一章,因为它依赖多个概念:

  • 节点和依赖 —— 第 2 章第 3 章
  • 批处理和迭代 —— 第 11 章引入了 foreachloop;streaming 变体在此基础上增加了 channel。
  • 测试 —— 第 18 章展示了如何对流式 graph 结果做断言。
  • 可观测性 —— 第 19 章涵盖了记录背压事件的监听器。

本附录将 streaming 的心智模型集中在一处,方便你在需要增量数据流时随时查阅。