第08章 - 失败策略与调度
方案驱动的数据处理框架,其价值不仅在于”能跑通”,更在于”跑不通时如何优雅地失败”。OpenGIS DAF 将方案建模为有向无环图(DAG),通过调度引擎决定执行顺序与并发度,并通过失败策略控制错误传播。本章结合演示方案 04,深入讲解失败策略、跳过传播、重试与超时机制。
8.1 调度引擎:从方案到执行
8.1.1 DAG 构建
方案中的每个处理项(Item)通过 inputs 绑定声明依赖关系:
external输入:依赖外部数据源,不构成 item 间依赖;upstream输入:sourceId指向另一个 item 的 id,构成一条有向边。
调度引擎读取全部 item 的绑定后,构建一张有向无环图。若存在环(A 依赖 B、B 又依赖 A),方案校验阶段会直接拒绝——这是内置 21 条校验规则之一。
8.1.2 Kahn 拓扑排序
调度引擎使用 Kahn 拓扑排序 确定执行顺序:不断取出”入度为零”(所有上游已完成)的节点执行,直到所有节点执行完毕。这保证了:
- 每个 item 执行时,其上游输出必然已就绪;
- 无依赖关系的 item 天然具备并行潜力。
8.1.3 串行与并行调度
PlanExecutionPolicy.maxParallelism(默认 4)控制最大并发度。调度引擎会尽量并行执行互不依赖的 item。例如演示方案 01 中,transform-schools 与 transform-boundary 互不依赖,可并行执行;而 buffer-schools 依赖 transform-schools,必须等其完成。
8.2 失败策略:错误如何传播
8.2.1 两种失败策略
PlanExecutionPolicy.failurePolicy 决定当一个 item 失败时,其余 item 如何处理:
| 策略 | 行为 | 适用场景 |
|---|---|---|
stopOnAny(默认) |
任一 item 失败即停止整个方案,未执行的 item 全部取消 | 分析链必须完整,缺一步结果无意义 |
continueIndependent |
失败 item 标记失败,但无依赖的独立 item 继续执行;依赖失败 item 的 item 被跳过 | 批量质检、多源处理,希望尽量多产出 |
8.2.2 跳过传播
在 continueIndependent 模式下,错误沿依赖链传播:
- 失败 item 本身 → 状态
failed; - 依赖失败 item 的 item → 状态
skipped(其上游输出不可用,无法执行); - 与失败 item 无依赖关系的 item → 正常执行。
8.3 演示方案 04:失败策略实战
演示方案 04-failure-policy.json 构造了三种典型 item,用于观察失败策略与跳过传播。
8.3.1 方案 JSON
{
"id": "demo-04-failure-policy",
"name": "失败策略演示(失败项 + 独立成功项 + 上游跳过项)",
"version": "1.0.0",
"group": "demo",
"items": [
{
"id": "step-broken-input",
"operatorId": "buffer",
"inputs": {
"source": {
"type": "external",
"sourceId": "data/missing-file.geojson"
}
},
"parameters": {
"distance": 500
},
"output": {
"adapterType": "geojson",
"targetPath": "output/demo04/should-not-exist.geojson"
}
},
{
"id": "step-independent-ok",
"operatorId": "field_calculator",
"inputs": {
"source": {
"type": "external",
"sourceId": "data/schools.geojson"
}
},
"parameters": {
"target_field": "note",
"expression": "\"独立演示项\"",
"field_type": "String"
},
"output": {
"adapterType": "console"
}
},
{
"id": "step-skipped-dependent",
"operatorId": "clip",
"inputs": {
"source": {
"type": "upstream",
"sourceId": "step-broken-input"
},
"clip": {
"type": "external",
"sourceId": "data/city_boundary.geojson"
}
},
"output": {
"adapterType": "geojson",
"targetPath": "output/demo04/also-should-not-exist.geojson"
}
}
],
"executionPolicy": {
"failurePolicy": "continueIndependent"
}
}
8.3.2 三个 item 的预期行为
| item | 算子 | 输入 | 预期状态 |
|---|---|---|---|
step-broken-input |
buffer | data/missing-file.geojson(不存在) |
failed |
step-independent-ok |
field_calculator | data/schools.geojson(存在,无依赖) |
succeeded |
step-skipped-dependent |
clip | source 指向 step-broken-input |
skipped |
8.3.3 运行与退出码
daf run --plan plans/04-failure-policy.json
预期输出:
step-broken-input因输入文件不存在而失败(数据源解析失败,属调度层错误);step-independent-ok与失败项无依赖,在continueIndependent策略下正常执行并输出到控制台;step-skipped-dependent因上游step-broken-input失败而被跳过;- CLI 退出码为 1(只要有失败或跳过,CLI 即返回 1,便于 CI 判断)。
对比实验:若把
failurePolicy改为默认的stopOnAny,则step-broken-input失败后,step-independent-ok也不会执行——整个方案立即停止。
8.4 重试与超时
8.4.1 Item 级执行策略
每个 item 可通过 executionPolicy 单独配置重试与超时:
{
"id": "buffer-schools",
"operatorId": "buffer",
"executionPolicy": {
"timeout": "00:05:00",
"maxRetries": 1,
"retryInterval": "00:00:02",
"exponentialBackoff": true
}
}
| 字段 | 默认值 | 说明 |
|---|---|---|
timeout |
"00:30:00" |
单次执行超时上限(ISO 8601 时长格式) |
maxRetries |
0 |
失败后的最大重试次数 |
retryInterval |
"00:00:05" |
重试间隔 |
exponentialBackoff |
true |
是否指数退避(每次重试间隔翻倍) |
8.4.2 哪些错误会重试
重要:maxRetries 仅对执行期错误(ERR_RT_*,如算子运行抛异常、超时)生效。输入/数据源解析失败(SCH_*,如文件不存在、上游不可用)属于调度层错误,在调度阶段直接失败,不参与重试。
这正是演示方案 04 中 step-broken-input 立即失败、而非反复重试的原因——文件不存在是确定性的,重试无意义。
8.5 取消与中断
daf run 支持 Ctrl+C 取消。取消时,调度引擎会停止调度新的 item,并尝试中断正在执行的 item,已完成的中间结果(标记了 isIntermediate)不会写入最终输出。
8.6 本章小结
- 方案被建模为 DAG,调度引擎用 Kahn 拓扑排序 确定执行顺序;
failurePolicy控制错误传播:stopOnAny立即停止,continueIndependent让独立项继续、依赖项跳过;- 错误沿依赖链传播:失败 → 依赖项被跳过;
maxRetries只对执行期错误(ERR_RT_*)生效,输入/数据源错误(SCH_*)直接失败不重试;- CLI 退出码 0 表示全部成功,1 表示存在失败或跳过,便于 CI 集成。