BLOGE
0.9.8-RC1在线版 · 事实校验 2026-09-15 · English
第 11 章 —— 批处理与迭代
承诺: 读完本章后,你将知道如何使用
foreach处理集合中的每一项,使用loop重复工作直到条件满足,并在单个 graph 中组合这两种构造——所有这些都无需在 operator 代码中编写任何for/while循环。
学习目标
- 编写一个
foreach块,对集合进行扇出——默认并行或顺序执行——并在主体中引用当前项和索引。 - 编写一个包含
until、max_iterations、delay和carry的loop块,表达有界迭代并在轮次之间转发状态。 - 使用三个 loop 特有的隐式变量 ——
carry.<path>、prev.<nodeId>.<path>和loopIteration—— 并解释每个变量在运行时如何解析。 - 在下游节点中使用
depends_on消费foreach或loop的聚合输出。 - 识别
foreach和loop何时可以在同一个 graph 中组合——以及何时stream foreach/stream loop是更好的选择。
前置条件
- 第 10 章 —— 用子图复用 —— 你需要理解
foreach和loop的主体会被编译为嵌套子图。
源示例
| 文件 | 展示内容 |
|---|---|
batch-order-processing.bloge | 最简单的并行 foreach —— 获取订单、验证、处理 |
batch-order-parallel.bloge | 并行 foreach,带有下游 summarize 节点消费聚合输出 |
sequential-transfer.bloge | foreach … sequential —— 逐一执行的银行转账 |
status-polling.bloge | 基本的 loop,带 until 和 loopIteration —— 轮询直到 READY |
cursor-pagination.bloge | 带多字段 carry 的 loop —— API 游标分页 |
retry-with-backoff.bloge | 基于 loop 的自定义重试,带 carry 和 loopIteration |
logistics-batch-dispatch.bloge | 同一个 graph 中的 foreach + loop —— 调度包裹,然后轮询状态 |
basic.bloge (foreach) | 一致性测试 fixture:最小的 foreach |
sequential-index.bloge | 一致性测试 fixture:带索引绑定的 foreach … sequential |
with-carry.bloge | 一致性测试 fixture:带 carry、prev 和 loopIteration 的 loop |
为什么这很重要
到目前为止,你的 graph 中每个节点最多运行一次。真实系统很少这么简单:
- 一个电商结账流程必须验证每一条订单行——而不是只验证一条。
- 一个数据同步任务必须翻页遍历 API 的所有页面,并向前传递游标。
- 一个部署流水线必须轮询外部系统,直到状态变为 READY。
在过程式代码中,你会使用 for 或 while。在 BLOGE 中,等价的构造是 foreach(处理集合中的每一项)和 loop(重复直到条件成立)。两者都是超级节点(hyper-node):它们在外层 graph 中表现为单个节点,但内部会编译为一个嵌套子图,引擎会像对待其他 graph 一样对其进行调度、重试和观测。
核心好处是:迭代语义——并行度、顺序、有界重试、carry 状态——都声明在 DSL 中,而不是埋在 operator 代码里。引擎拥有循环;你的 operator 只需处理单个项或单次迭代。
心智模型
foreach —— 对集合进行扇出
把 foreach 想象成一台冲压机:引擎拿到集合,为每个元素冲出一份主体子图的副本,运行它们(默认并行,或顺序执行),然后将结果收集为聚合输出。
loop —— 重复直到完成
把 loop 想象成一个带计数器的闸门:引擎重新执行嵌套子图,在每轮结束后评估 until,当谓词为 true 或 达到 max_iterations 时停止。在迭代之间,引擎等待 delay 指定的时间,并将 carry 状态传入下一轮。
第一个可运行示例
下面是
batch-order-processing.bloge
—— 最简单的并行 foreach:
graph batchOrderProcessing {
node fetchOrders : OrderFetcher {
input {
customerId = ctx.customerId
}
}
foreach processOrders : (order, idx) in fetchOrders.output.orders {
node validate : OrderValidator {
input {
order = order
index = idx
}
}
node process : OrderProcessor {
depends_on = [validate]
input {
validated = validate.output
}
}
}
}
逐步解读:
fetchOrders运行并产出output.orders—— 一个列表。- 引擎进入
foreach processOrders。对于列表中的每一个元素,它创建一份主体子图的副 本,将order绑定到当前元素,idx绑定到其从 0 开始的位置。 - 默认情况下,所有副本并行运行。
- 在每个副本中,
validate先运行;process依赖于它。 - 当所有副本完成后,
processOrders本身以聚合输出完成。
关键认识: 你从不编写循环。引擎负责扇出、调度和收集。你的 operator ——
OrderValidator、OrderProcessor—— 只处理单个项。
完整故事 A:foreach 交付一张批处理 receipt
orderBatch 加载三笔订单,对每一项运行相同的 processOrder 能力,再把一组
可关联结果交给 generateReport。结果不是“后台发生了三次任务”,而是一张
父图可以检查的批处理 receipt。
分开观察输入、单项结果与汇总
| 层级 | 示例观察 | 失败时要问什么 |
|---|---|---|
| collection | 输入包含 O-1、O-2、O-3 | dispatch 前是否漏掉某项? |
| item | index 1 对应 O-2 | 这一项执行了还是失败了? |
| aggregate | totalProcessed=3 | 父图是否收到完整 receipt? |
companion Scenario 断言 generateReport 已执行且 totalProcessed=3。它证明
三项 fixture 在这次受控运行中到达报告;它不证明无限并行、所有失败策略下
的顺序,也不证明外部 effect 安全。
完整故事 B:loop 交付一张退出 receipt
轮询 loop 回答的是另一个问题。submitJob 只运行一次;loop 复用
jobId=J-42 检查状态,直到看到 READY 或触及 max_iterations。
loop 必须同时拥有记忆和停止规则
| receipt 字段 | 示例 | 父图为什么需要它 |
|---|---|---|
iterations | 3 | 区分立即成功与重复轮询 |
exit | until | 区分业务条件成立与保护上限 |
| 最终输出 | status=READY | 不暴露每轮工作台也能给 fetchResult 供值 |
上一轮输出,以及 DSL 明确声明的 carry 字段,共同构成 loop 的记忆。本例的
jobId 始终来自父图输入,并没有被伪装成 carry 状态。until 是业务停止条件,
max_iterations 是安全停止条件。如果达到安全上限时仍为 PENDING,不能因为引擎正常结束就把轮询写成
成功。
foreach 把一个 collection 展开为多个兄弟执行;loop 让一个逻辑活动跨轮
演进。后面可以组合两者,但前提是两张 receipt 都能独立读懂。
拆开来看
foreach 语法
正式语法 (DSL 规范 §3.1):
ForEachDef = "foreach" IDENT ":" ForEachBinding "in" Expression
"sequential"? "{" body "}"
ForEachBinding = "(" IDENT ("," IDENT)? ")" // 带可选索引的元组
| IDENT // 裸项变量
| 部分 | 示例 | 含义 |
|---|---|---|
| ID | processOrders | 外层 graph 中超级节点的名称 |
| 绑定 | (order, idx) | order = 当前项,idx = 从 0 开始的索引。两者都是可选的——(order) 或裸 order 都有效。 |
in 表达式 | fetchOrders.output.orders | 要迭代的集合 |
sequential | (关键字,可选) | 逐一处理项,而非并行 |
| 主体 | { node … } | 编译为 <foreachId>__subgraph__ 的嵌套子图 |
foreach 内部的隐式变量
| 变量 | AST 类型 | 运行时上下文键 | 含义 |
|---|---|---|---|
<itemVar>(如 order) | ItemPath | __item__ | 当前集合元素 |
<itemVar>.<field>(如 order.amount) | ItemPath | __item__ | 当前元素上的字段 |
<indexVar>(如 idx) | ItemIndex | __itemIndex__ | 从 0 开始的整数位置 |
(参见 DSL 规范规则 7–8。)
并行 vs. 顺序
在
sequential-transfer.bloge
中,银行转账必须按顺序处理,以便每笔转账能看到前一笔更新后的余额:
foreach processTransfers : (transfer, idx) in fetchTransfers.output.transfers sequential {
node riskCheck : RiskCheckOperator {
input {
amount = transfer.amount
fromAccount = transfer.fromAccount
toAccount = transfer.toAccount
index = idx
}
}
node executeTransfer : TransferExecutionOperator {
depends_on = [riskCheck]
input {
transfer = transfer
riskResult = riskCheck.output
}
}
node recordLedger : LedgerRecordOperator {
depends_on = [executeTransfer]
input {
transfer = transfer
execution = executeTransfer.output
index = idx
}
}
}
唯一的语法区别是在 in 表达式后面加上 sequential 关键字。运行时引擎按列表顺序逐一执行每个项。