Skip to main content

BLOGE 0.9.8-RC1 在线版 · 事实校验 2026-09-15 · English

第 11 章 —— 批处理与迭代

承诺: 读完本章后,你将知道如何使用 foreach 处理集合中的每一项,使用 loop 重复工作直到条件满足,并在单个 graph 中组合这两种构造——所有这些都无需在 operator 代码中编写任何 for/while 循环。


学习目标

  1. 编写一个 foreach 块,对集合进行扇出——默认并行或顺序执行——并在主体中引用当前项和索引。
  2. 编写一个包含 untilmax_iterationsdelaycarryloop 块,表达有界迭代并在轮次之间转发状态。
  3. 使用三个 loop 特有的隐式变量 —— carry.<path>prev.<nodeId>.<path>loopIteration —— 并解释每个变量在运行时如何解析。
  4. 在下游节点中使用 depends_on 消费 foreachloop 的聚合输出。
  5. 识别 foreachloop 何时可以在同一个 graph 中组合——以及何时 stream foreach / stream loop 是更好的选择。

前置条件

源示例

文件展示内容
batch-order-processing.bloge最简单的并行 foreach —— 获取订单、验证、处理
batch-order-parallel.bloge并行 foreach,带有下游 summarize 节点消费聚合输出
sequential-transfer.blogeforeach … sequential —— 逐一执行的银行转账
status-polling.bloge基本的 loop,带 untilloopIteration —— 轮询直到 READY
cursor-pagination.bloge带多字段 carryloop —— API 游标分页
retry-with-backoff.bloge基于 loop 的自定义重试,带 carryloopIteration
logistics-batch-dispatch.bloge同一个 graph 中的 foreach + loop —— 调度包裹,然后轮询状态
basic.bloge (foreach)一致性测试 fixture:最小的 foreach
sequential-index.bloge一致性测试 fixture:带索引绑定的 foreach … sequential
with-carry.bloge一致性测试 fixture:带 carryprevloopIterationloop

为什么这很重要

到目前为止,你的 graph 中每个节点最多运行一次。真实系统很少这么简单:

  • 一个电商结账流程必须验证每一条订单行——而不是只验证一条。
  • 一个数据同步任务必须翻页遍历 API 的所有页面,并向前传递游标。
  • 一个部署流水线必须轮询外部系统,直到状态变为 READY。

在过程式代码中,你会使用 forwhile。在 BLOGE 中,等价的构造是 foreach(处理集合中的每一项)和 loop(重复直到条件成立)。两者都是超级节点(hyper-node):它们在外层 graph 中表现为单个节点,但内部会编译为一个嵌套子图,引擎会像对待其他 graph 一样对其进行调度、重试和观测。

核心好处是:迭代语义——并行度、顺序、有界重试、carry 状态——都声明在 DSL 中,而不是埋在 operator 代码里。引擎拥有循环;你的 operator 只需处理单个项或单次迭代。


心智模型

foreach —— 对集合进行扇出

Diagram: 11-batch-and-iteration figure 1

foreach 想象成一台冲压机:引擎拿到集合,为每个元素冲出一份主体子图的副本,运行它们(默认并行,或顺序执行),然后将结果收集为聚合输出。

loop —— 重复直到完成

Diagram: 11-batch-and-iteration figure 2

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
}
}
}
}

逐步解读:

  1. fetchOrders 运行并产出 output.orders —— 一个列表。
  2. 引擎进入 foreach processOrders。对于列表中的每一个元素,它创建一份主体子图的副本,将 order 绑定到当前元素,idx 绑定到其从 0 开始的位置。
  3. 默认情况下,所有副本并行运行。
  4. 在每个副本中,validate 先运行;process 依赖于它。
  5. 当所有副本完成后,processOrders 本身以聚合输出完成。

关键认识: 你从不编写循环。引擎负责扇出、调度和收集。你的 operator —— OrderValidatorOrderProcessor —— 只处理单个项。


完整故事 A:foreach 交付一张批处理 receipt

orderBatch 加载三笔订单,对每一项运行相同的 processOrder 能力,再把一组 可关联结果交给 generateReport。结果不是“后台发生了三次任务”,而是一张 父图可以检查的批处理 receipt。

图:foreach fan-out 与汇总 receipt

分开观察输入、单项结果与汇总

层级示例观察失败时要问什么
collection输入包含 O-1O-2O-3dispatch 前是否漏掉某项?
itemindex 1 对应 O-2这一项执行了还是失败了?
aggregatetotalProcessed=3父图是否收到完整 receipt?

companion Scenario 断言 generateReport 已执行且 totalProcessed=3。它证明 三项 fixture 在这次受控运行中到达报告;它不证明无限并行、所有失败策略下 的顺序,也不证明外部 effect 安全。


完整故事 B:loop 交付一张退出 receipt

轮询 loop 回答的是另一个问题。submitJob 只运行一次;loop 复用 jobId=J-42 检查状态,直到看到 READY 或触及 max_iterations

图:loop 迭代与退出 receipt

loop 必须同时拥有记忆和停止规则

receipt 字段示例父图为什么需要它
iterations3区分立即成功与重复轮询
exituntil区分业务条件成立与保护上限
最终输出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 // 裸项变量
部分示例含义
IDprocessOrders外层 graph 中超级节点的名称
绑定(order, idx)order = 当前项,idx = 从 0 开始的索引。两者都是可选的——(order) 或裸 order 都有效。
in 表达式fetchOrders.output.orders要迭代的集合
sequential(关键字,可选)逐一处理项,而非并行
主体{ node … }编译为 <foreachId>__subgraph__ 的嵌套子图

foreach 内部的隐式变量

变量AST 类型运行时上下文键含义
<itemVar>(如 orderItemPath__item__当前集合元素
<itemVar>.<field>(如 order.amountItemPath__item__当前元素上的字段
<indexVar>(如 idxItemIndex__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 关键字。运行时引擎按列表顺序逐一执行每个项。

消费 foreach 的下游输出

一个 foreach 超级节点会暴露聚合输出。任何下游节点都可以通过 depends_on 依赖该 foreach 并读取它。摘自 batch-order-parallel.bloge

node summarize : BatchSummaryOperator {
depends_on = [processOrders]
input {
results = processOrders.output
}
}

loop 语法

LoopDef = "loop" IDENT "{" LoopBody "}"

LoopBody = ( "max_iterations" "=" NUMBER
| "delay" "=" DURATION
| DependsOn
| Member
| CarryBlock
| UntilClause )*

CarryBlock = "carry" "{" ( IDENT ":" Expression ","? )* "}"
UntilClause = "until" Expression
部分示例含义
IDpollStatus外层 graph 中超级节点的名称
max_iterations20硬性上限——即使 until 从未触发,引擎也会停止
delay2s迭代之间的暂停时间
untilcheckStatus.output.status == "READY"提前退出谓词——在每次迭代后评估
carry { … }cursor: fetchPage.output.nextCursor转发到下一次迭代的命名状态
主体node …编译为 <loopId>__subgraph__ 的嵌套子图

loop 内部的隐式变量

变量AST 类型运行时上下文键含义
carry.<path>LoopCarryPath__carry__来自上一次迭代(或初始输入)的状态 map
prev.<nodeId>.<path>LoopPrevPath__prev__上一次迭代的节点输出
loopIterationLoopIterationRef__loopIteration__从 0 开始的迭代计数器

(参见 DSL 规范规则 9–11。)

基本 loop —— 状态轮询

摘自 status-polling.bloge

loop pollStatus {
max_iterations = 20
delay = 2s
depends_on = [submitJob]
node checkStatus : StatusCheckerOperator {
input {
jobId = submitJob.output.jobId
iteration = loopIteration
}
}
until checkStatus.output.status == "READY"
}

引擎的执行流程:

  1. 运行 checkStatus(第 0 次迭代)。
  2. 评估 until。如果 status == "READY" → 退出。
  3. 否则,等待 2 秒,递增 loopIteration,然后重复。
  4. 经过 20 次迭代仍未匹配后,循环停止 —— max_iterations 是安全网。

带 carry 的 loop —— 游标分页

摘自 cursor-pagination.bloge

loop fetchAllPages {
max_iterations = 100
delay = 500ms
depends_on = [initPagination]

node fetchPage : PageFetcherOperator {
input {
endpoint = ctx.endpoint
pageSize = ctx.pageSize
cursor = carry.cursor
}
}
node transformPage : PageTransformerOperator {
depends_on = [fetchPage]
input {
records = fetchPage.output.records
currentTotal = carry.totalRecords
}
}

carry {
cursor: fetchPage.output.nextCursor
totalRecords: transformPage.output.runningTotal
}
until fetchPage.output.hasMore == false
}

要点:

  • carry.cursorcarry.totalRecords 在主体中被读取——它们来自上一次迭代的 carry { … } 块。
  • 第一次迭代时,carry 字段未设置(null),因此 operator 应当处理初始情况。
  • carry 块定义了在每次迭代结束时写入的状态,供下一次迭代读取。

消费 loop 的下游输出

loop 完成后,其输出是最后一次迭代的节点输出,以节点 ID 为键:

node finalizeData : DataFinalizerOperator {
depends_on = [fetchAllPages]
input {
totalRecords = fetchAllPages.output.transformPage.runningTotal
lastCursor = fetchAllPages.output.fetchPage.nextCursor
}
}

在同一个 graph 中组合 foreachloop

logistics-batch-dispatch.bloge 展示了两种构造协同工作:

  1. 一个 foreach 对包裹进行扇出(路线规划 → 调度)。
  2. 一个 loop 在 foreach 完成之后轮询聚合的配送状态,通过 depends_on = [assignRoutes] 连接。
  3. 最终的 dispatchReport 节点依赖于 loop。

这展示了核心组合规则:foreachloop 是超级节点——你像连接其他节点一样用 depends_on 连接它们。


关于流式变体的说明

DSL 还支持 stream foreachstream loop。它们会在每个项或每次迭代的结果完成时立即向下游发送,而不是等待整个批次完成。例如,摘自 streaming-batch.bloge

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

流式处理是一个超出本章范围的高级主题。关键要知道的是它的存在,以及 buffer = N 控制流式边的环形缓冲区容量。


限流与单项容错

并行 foreach 默认是"一次性把所有 item 全部启动"。在小列表上很快, 但对一个只能扛 50 路并发的下游服务一下子扇出 10,000 个 item,就是灾难。 有两个 foreach 属性可以解决,且不用动 body:

foreach order in orders {
batch_size = 25 // 同一时刻最多 25 个 item 在跑
on_item_failure = continue // 一个坏 item 不会拖垮整批

node process : ProcessOrderOperator {
input { order = item }
}
}

batch_size

取值含义
省略不限流 —— 引擎为每个 item 启动一个任务
batch_size = N(N ≥ 1)引擎同一时刻最多处理 N 个 item;一个完成就启动下一个

batch_size 是局部限流。全局封顶请看下面的 maxGlobalForeachBatches

on_item_failure

取值单个 item 失败时的行为
abort_all(默认)把异常向上抛;整个 foreach 失败 —— 经典 fail-fast 语义
abort_batch把当前在跑的 item 跑完,然后 foreach 标记失败;后续 item 不再启动
continue把失败记进聚合输出,继续处理剩下的;foreach 自身成功完成

设为 on_item_failure = continue 时,聚合输出保留输入顺序,并为失败 item 留下哨兵记录,下游节点可以按 item.status 分支:

node summarize : SummarizeBatchOperator {
depends_on = [process]
input {
results = process.output // List,既有成功项也有失败项
}
}

全局封顶 —— maxGlobalForeachBatchesdefaultForeachProtection

batch_size 是单个 foreach 上的设置。在繁忙的服务里你还想要一个 系统级封顶,防止某个热点租户把整个服务挤死。这两个引擎级开关对所有 图都生效:

GraphEngine engine = GraphEngine.builder()
.maxGlobalForeachBatches(200) // 跨所有图的硬上限
.defaultForeachProtection( // 对 .bloge 里未声明 batch_size
ForeachProtection.builder() // 的 foreach 的兜底
.batchSize(50)
.onItemFailure(OnItemFailure.ABORT_BATCH)
.build())
.build();
  • maxGlobalForeachBatches:集群级上限。触顶后新批次排队等已有的跑完。
  • defaultForeachProtection:对 .bloge 中未指定 batch_size / on_item_failureforeach 应用的兜底配置。生产环境务必配置,以免 漏配某个开关就把集群打爆。

小测验: on_item_failure = continue 什么时候应该配人工告警? 提示:50 个 item 失败了 49 个,批次还显示 "OK",你就掩盖了一次真实故障。 continue 必须搭配一个下游节点检查失败率。


常见陷阱

❌ 忘记 carry 字段在第一次迭代时为 null

// 错误 —— 在第 0 次迭代时 cursor 将为 null;operator 可能会崩溃
loop fetchPages {
max_iterations = 50
node fetch : PageFetcher {
input {
cursor = carry.cursor // 第一次迭代时为 null!
}
}
carry { cursor: fetch.output.nextCursor }
until fetch.output.done == true
}

第一次迭代时,没有上一轮的 carry —— 每个 carry.<field> 都解析为 null。你的 operator 必须处理这种情况:

// 在 PageFetcher.execute(...) 中
String cursor = (String) input.get("cursor");
if (cursor == null) {
// 第一页 —— 从头开始
}

这不是 DSL 的 bug。它和循环变量在第一次迭代前未初始化是同一个道理。设计你的 operator 时,将 null carry 值视为"从头开始"。

❌ 在 loop 上省略 max_iterations

max_iterations硬性安全网。没有它,一个不稳定的 until 条件可能导致循环无限运行。即使你预期 until 会远在上限之前触发,也要始终声明一个合理的上限。


常见故障

loop 缺少 max_iterations

loop pollStatus {
until = checkStatus.output.status == "READY"
delay = 5s
// No max_iterations — the compiler rejects this
body {
node checkStatus : StatusCheckerOperator { … }
}
}

解析器会直接拒绝这种写法:

GraphDefinitionException: Loop 'pollStatus' is missing required 'max_iterations'

修复: 始终声明 max_iterations。选择一个能够覆盖现实最坏情况、并留有余量的上限。


引导式重写

从一致性测试 fixture with-carry.bloge 开始:

graph g {
loop processBatch {
max_iterations = 100
delay = 5s
node fetchBatch : BatchFetcher {
input {
cursor = carry.cursor
prevResult = prev.fetchBatch.result
iteration = loopIteration
}
}
carry { cursor: fetchBatch.output.nextCursor }
until fetchBatch.output.done == true
}
}

思考以下问题:

  1. 这里使用了哪三个隐式变量? carry.cursor —— 来自上一次迭代的 carry 状态。prev.fetchBatch.result —— 上一次迭代 fetchBatch 输出的 result 字段。loopIteration —— 从 0 开始的迭代计数器。)

  2. 在第 0 次迭代时,carry.cursor 的值是什么? (Null —— 还没有上一轮的 carry。operator 必须处理这种情况。)

  3. 在第 0 次迭代时,prev.fetchBatch.result 的值是什么? (Null —— 没有上一次迭代。与 carry 相同的规则。)

  4. 如果 fetchBatch.output.done 永远不变为 true,会怎样? (循环运行 100 次 —— max_iterations 是硬性停止条件 —— 然后退出。)

现在自己扩展这个 graph:

  • 在 loop 主体中添加一个 node transformBatch,使其依赖 fetchBatch,对原始记录进行转换。
  • 更新 carry 块,增加 totalRecords: transformBatch.output.runningTotal
  • 在 loop 外部添加一个下游 node finalize,消费 processBatch.output.transformBatch.runningTotal

将你的结果与 cursor-pagination.bloge 进行对比——它遵循相同的模式。


脑力检查

  1. foreach 的默认执行模式是并行还是顺序? (并行。添加 sequential 关键字可切换为逐一执行。)

  2. foreach 主体中有哪两个隐式变量可用? (项变量——如 order——解析为 ItemPath;以及可选的索引变量——如 idx——解析为 ItemIndex。)

  3. 列出 loop 主体中可用的三个隐式变量。 carry.<path> —— LoopCarryPathprev.<nodeId>.<path> —— LoopPrevPathloopIteration —— LoopIterationRef。)

  4. until 永远不为 true 时会怎样? (循环运行到 max_iterations 为止,然后退出。)

  5. 下游节点如何消费 foreachloop 的输出? (通过声明 depends_on = [<foreachOrLoopId>] 并读取 <id>.output。)

  6. 基于 loop 的重试和节点级 retry 有什么区别? (基于 loop 的重试让你在每次尝试中运行任意子图、携带自定义状态并计算动态退避。节点级 retry 是引擎自动应用于单个 operator 调用的声明式策略。)

  7. 设计题: 你有 1 000 个订单要处理。一个并行 foreach 会立刻扇出到 1 000 个虚拟线程,但下游支付网关只允许 50 个并发请求。你会如何处理?(引擎本身没有内建的 foreach 并发限流。常见做法是:先把订单切成每批 50 个的分组、在 operator 内部做限流或信号量控制,或者如果顺序本来就重要,就改用 foreach sequential。真正的设计取舍在于你更看重吞吐量还是顺序约束。)


练习

  1. 并行 foreach —— 批量验证 打开 batch-order-parallel.bloge

    • foreach 主体中添加第三个节点 notifyCustomer,使其依赖 deductStock
    • 添加一个 input 绑定,读取 order.orderIddeductStock.output.remaining
    • 预测:添加这个节点是否会改变跨项的并行度? (不会——每个项的子图仍然是独立的。新节点只在每个项的副本内增加了一个步骤。)
  2. 顺序 foreach —— 顺序很重要 打开 sequential-transfer.bloge

    • 删除 sequential 关键字。
    • 预测:会发生什么变化? (所有转账现在并行执行——风险检查和余额扣减交织进行。这对于顺序很重要的银行转账来说是错误的。)
  3. 带 carry 的 loop —— 构建分页器 从头编写一个新 graph:

    • node initSearch 提供初始查询。
    • loop fetchPages,设置 max_iterations = 50delay = 1s
    • loop 内部:node fetchPage 读取 carry.pageToken
    • carry { pageToken: fetchPage.output.nextPageToken }
    • until fetchPage.output.hasMore == false
    • loop 之后:node aggregate 消费 fetchPages.output
    • cursor-pagination.bloge 进行对比。
  4. 组合 foreach + loop 打开 logistics-batch-dispatch.bloge

    • 追踪数据流:哪个节点的输出流入 foreach?哪个 depends_onforeach 输出连接到 loop?哪个将 loop 连接到最终报告?
    • 在 foreach 内部的 planRoute 上添加 retry = { attempts: 2, backoff: 500ms, strategy: fixed }。这会影响 loop 吗? (不会——retry 是按节点的,在每个 foreach 项的子图内部。)

实验验收卡

  • 预期与观察: foreach 交付逐项 receipt,loop 交付退出 receipt。
  • 失败与恢复: 让 loop 永不满足 until;恢复上限或 carry 更新。
  • 证明边界: 证明批处理与迭代合同不同,不证明无限流自动收敛。
  • 练习合同: 三项批次和一个 loop;只改退出条件;交付两类 receipt;关联完整且退出原因唯一即停止。

回顾

  • foreach 对集合进行扇出。每个项获得主体子图的一份副本。默认并行;添加 sequential 可按顺序执行。
  • foreach 内部,项变量(如 order)和可选的索引变量(如 idx)是隐式的——你在绑定中声明它们,在 input 表达式中使用它们。
  • loop 重复执行一个嵌套子图,直到 until 评估为 true 或达到 max_iterationsdelay 在轮次之间暂停。
  • loop 内部,三个隐式变量提供迭代上下文:carry.<path>(转发状态)、prev.<nodeId>.<path>(上一次迭代的输出)和 loopIteration(从 0 开始的计数器)。
  • 两种构造都编译为超级节点 —— 它们在外层 graph 中表现为单个节点,像其他节点一样通过 depends_on 连接。
  • carry 块定义在每次迭代结束时写入的状态。在第一次迭代时,carry 字段为 null —— operator 必须处理这种情况。
  • 流式变体(stream foreachstream loop)在结果完成时即向下游发送,而不是等待整个批次完成。

下一步

第 12 章——等待世界的回应中,你将学习 BLOGE graph 如何暂停并等待外部事件——定时器、人工审批和 webhook 回调——而不阻塞线程。


参考链接

Coding Agent: Open the versioned task guide.