WRITING / 2026.09.23

Rill Flow 专题(三):复杂流程如何表达——数据映射、分支与并行汇合

以视频处理流程为例,解析 Rill Flow 的输入输出映射、上下文作用域、条件分支与并行汇合,理解流程结构正确之外还需要维护哪些数据约束。

这是 Rill Flow 系列第三篇。第二篇已经串起一次执行的推进路径,本文转向流程定义:一张图怎样同时表达依赖、数据和分支。源码固定到公开仓库 5cded0b 版本,限定普通非流式输入、未启用关键路径提前完成的场景。视频链路是教学案例;配置经过语法检查,运行行为来自源码推演,未启动引擎执行。

依赖正确,为什么仍可能拿错结果

假设探测服务返回视频信息,分片服务生成三个片段,转码完成后合并。图上每条箭头都正确,合并却只收到一个文件。原因可能很简单:三个任务把产物都写到了同一个结果字段,后一次更新覆盖了前一次结果。调度条件回答了谁先谁后,却没有回答每份数据属于谁。

还有一种更隐蔽的错误:三个片段都已成功,输出列表却按完成顺序排列。最短的片段先完成,被拼在视频开头。引擎完成了任务协调,业务结果仍然错误。这里缺少的是稳定的分片身份和合并顺序,不能依靠执行速度推断。

因此,设计流程时要分别画出控制关系和数据关系。前者说明哪些任务必须等待;后者说明哪个输入来自哪个结果、字段在哪个作用域保存、汇总时怎样识别一项结果。只有两者都闭合,图才足以支撑业务。

本文给每个分片约定序号、源文件地址和处理参数版本,给转码产物约定序号与结果地址。这些是案例中的应用契约,不是 Rill Flow 强制提供的字段。名字可以变化,但身份、版本与结果之间的对应关系不能省略。

流程定义如何成为执行模型

DAGStringParser把描述文本解析为模型,FlowDAGValidator检查命名、节点引用和路径中的环等约束。两步的职责不同:YAML 能解析,说明文本符合格式;结构校验通过,才说明模型满足引擎检查的图约束。它们都不能代替执行器接口的业务校验。

例如,next 指向不存在的任务是结构问题;指向存在的合并任务,但没有传入文件列表,则是数据契约问题。再如,输入字段叫 segment,执行器却只接受 segment_url,这张图仍可能在远程调用时失败。不能把“编辑器保存成功”理解成完整接入验证。

同一份定义可以创建多个执行实例。视频甲和视频乙的上下文分别属于自己的执行,节点名称则用于在定义和实例中定位任务。排障时至少需要同时知道执行实例和任务位置;只拿到一个叫 transcode 的名字,无法判断它属于哪个视频、哪一组子任务。

结构也影响数据作用域。该版本校验器要求任务定义的名称全局不重复;foreach 各组从同一子任务定义展开,运行时生成不同的任务标识和分组上下文。不能因为使用同一份子任务定义,就把各组输入当作同一个全局对象。阅读配置时,缩进不仅是展示层级,也决定了任务归属。

本文只保留与问题相关的 YAML 片段,不提供一份完整可部署的视频工作流。真实接入还需要初始上下文、执行器地址及返回协议;这些缺失部分不能靠示例字段名自动补齐。需要完整结构时,可对照固定版本的官方并行异步示例,注意它处理的是数值而不是视频。

输入、输出与上下文如何连接

JSONPathInputOutputMapping在映射时组织 contextinputoutput 三个入口。输入映射为当前调用准备参数,输出映射把调用结果写回可供后继读取的位置。它们关联同一次执行,却承担不同职责,不能把三者当作一个随意读写的大对象。

假设某个分组的上下文包含下面的教学数据。地址仅作标识,不对应可访问的媒体:

{
  "segment": {"index": 1, "url": "segment-001.mp4"},
  "profileVersion": "h264-v2"
}

对应输入映射可以只挑出服务真正需要的字段:

inputMappings:
  - source: $.context.segment.url
    target: $.input.source_url
    tolerance: false
  - source: $.context.segment.index
    target: $.input.segment_index
    tolerance: false

按这两条规则推演,当前节点的输入如下。它是预期快照,不是引擎运行输出:

{"source_url": "segment-001.mp4", "segment_index": 1}

处理参数版本没有被自动复制到输入中。如果服务需要它,就必须增加映射。这也是使用显式映射的收益:可以检查每个依赖字段,而不是让执行器隐式依赖整个上下文的内部结构。

映射容错需要单独看待。该实现只有在规则明确设置 tolerance: false 时,才会把映射过程中捕获的异常继续抛出;容错时则记录并略过该规则。但这不意味着严格模式会验证所有必填字段:转换结果为 null 时,代码不会执行目标赋值,未必发生异常。必填、类型与取值范围仍要由业务校验说明。

任务的 tolerance 又是另一层配置,它影响任务失败后的处理。映射漏掉一个参数,与执行器运行失败后允许跳过,不是同一件事。把两者都写成“容错开启就继续”,会隐藏结果不完整的原因。

作用域可以用下图理解。图中路径是概念表示,不是存储键格式:

执行实例上下文
  └─ foreach 父节点
       ├─ 分组 0:segment → result
       ├─ 分组 1:segment → result
       └─ 分组 2:segment → result
                ↓ 父节点汇总
          合并节点的输入

子上下文让各分组能够使用相同的业务字段名,并不等于执行器可以无条件修改任意父级字段。选择字段时应先确定所属实例、所属组及消费节点,再决定映射位置。大文件放在外部对象存储、上下文只保存引用和校验信息,是本文建议的业务设计,可减少数据复制和后续排障负担。

条件分支如何改变后继运行条件

视频探测完成后,可能需要根据编码格式决定直接进入后续步骤,还是先做额外处理。SwitchTaskRunner对配置的条件逐项判断,记录需要运行和跳过的后继任务。这里使用的是该版本的 switch 节点,不把它与另一个 choice 实现混为一谈。

条件命中不天然意味着只选择一个分支。源码会维护已命中的后继集合,只有命中的规则设置了 break,才会影响后续条件的判断。默认条件则在普通条件处理之后统一判断;如果此前已有普通条件命中,就不再需要默认分支。是否互斥必须由条件与配置共同决定。

下面是阅读判断路径得到的语义整理,不是执行结果:

配置或输入情况 判断方向 写流程时要确认
多条普通条件命中,未中断 多个后继可被选中 是否真的允许多分支执行
某条命中且设置中断 后续条件不再正常命中 规则顺序是否符合业务优先级
普通条件均未命中 判断默认条件 是否定义了有效兜底路径
多条规则指向同一后继 任一命中可保留该后继 不把未命中的规则单独当成最终跳过决定

条件表达式异常也值得注意:这条实现路径会记录异常并返回不匹配。配置作者因此应准备输入样例,逐项检查条件,而不能期待所有表达式错误都在发布定义时被发现。默认分支最好表达一个明确业务结果,例如拒绝不支持的格式,而不是悄悄当作成功。

跳过状态还会影响汇合。假设两条分支中只应执行一条,未选中的分支被跳过,后继需要结合成功与跳过状态判断依赖。如果业务把“跳过”解释成“对应产物必然存在”,就会在汇合处再次失败。控制语义允许前进,并不负责替未执行的节点生成数据。

这也是分支节点与数据映射必须一起审查的原因。合并或结果封装节点应明确读取哪条已选择路径的字段;互斥分支可以约定统一输出结构,但统一结构需要两边主动维护,不能仅靠相同的后继名称保证。

foreach 如何展开任务并收集结果

ForeachTaskRunner根据 iterationMapping.collection 取得集合,为每一项创建子任务组和子上下文。构建子上下文时,它会复制处理后的公共输入,再放入当前迭代项;配置了索引名称时,还会写入分组索引。由此,每个分片既能取得公共参数,也能定位自己的输入。

父节点会记录各组状态。普通模式下,DAGWalkHelper在全部分组成功或跳过时把父节点判断为成功;如果仍有运行或待运行的组,就仍处于推进过程中。业务若要求所有分片都生成产物,必须额外限制跳过的使用,不能把父节点成功直接当成文件完整性证明。

空集合是一个典型边界。该 runner 在集合为空或没有子任务时,会直接将任务设为成功。对于“批量发送可选通知”,这可能合理;对于“至少包含一个片段的视频合并”,却可能意味着上游已经出错。因此,空集合是否允许,应在业务校验节点表达,不能等待并行结构替业务作决定。

父节点汇总时,可通过 sub_context 收集各组写入的字段。下面只展示输出映射片段;假设每组已将序号和地址组成 result 对象写入所属上下文:

outputMappings:
  - source: $.output.sub_context.[*].result
    target: $.context.transcoded_segments

相比仅收集 URL,保留完整结果对象更容易验证对应关系。但仍应由合并服务按业务序号排序、排除重复并检查缺片,不能把上下文枚举顺序当成约定好的播放顺序。分片产物还应带上可核对的处理版本,防止一次重试拿到旧参数生成的文件。

并行数量也不等于集合长度。该实现存在 synchronization 控制入口,其生效需要满足对应条件与配置要求。本文只解释分组和汇总,运行治理将在第五篇展开;不要从一条并发配置推断整个集群的最大在途任务数已经受到同样约束。

把流程正确性落到数据契约

为这条视频链路设计验收时,应从最终结果倒推每一层责任。先明确合并需要多少片、顺序如何确定、哪些编码属性必须兼容,再检查分片输出和转码输出能否提供这些信息。最后才检查映射是否把信息送到了正确位置。

一个足够小的契约可以约定:分片列表中的序号唯一;每项都有源文件标识;每个有效结果对应同一输入与处理版本;合并前结果数量和序号集合与预期一致。计数相同仍可能存在重复片段,因此数量检查必须和身份检查配合。这些规则可以由普通业务函数执行,不必为了表达它们修改引擎。

例如,分片阶段声明期望序号为零、一、二,汇总拿到零、一、一,即使三个转码节点都返回了成功,合并也必须拒绝。若拿到二、零、一,则可以在验证身份之后排序,不必重做计算。前一种是完整性错误,后一种只是表现顺序不同;把两者都处理成工作流失败并重试,会浪费已有结果,也使真正的问题更难定位。

还要规定字段的缺省含义。没有结果字段、字段值为 null、结果列表为空,可能分别表示尚未赋值、明确没有产物和没有待处理项目。业务若允许这些状态,就应写出各自的处理分支;如果不允许,应尽早产生可解释的失败。依赖映射层静默跳过缺失值,只会把错误推迟到更难定位的后继节点。

输出契约最好随定义一起评审:每个节点列出必填输入、成功输出、允许缺失的字段和失败后是否留有部分结果。这样既能检查新分支是否满足合并需求,也能判断容错跳过是否合理。配置评审关注的是一个节点对后继作出的承诺,而不只是某条路径能否在图上连通。

配置评审至少应覆盖五类输入:正常三片、空集合、缺少必填字段、结果乱序和某片失败。对分支再增加多条件命中与全部不命中两类情况。本文只依据源码整理预期行为,并对示例进行语法检查;要证明部署后的实际行为,还需在隔离环境使用对应执行器运行这些案例。

后续变更也应沿着同一契约审查。某个节点把 result_url 改成嵌套对象,即使输出更丰富,也可能破坏旧映射;新增一个分支,则可能让原本总会存在的字段变为可选。定义和执行器协议的版本要一起管理,不能只比较流程图新增了几条线。

本篇可按映射实现分支实现集合展开的顺序回读源码,并对照映射测试源码并行测试源码中的输入与断言。官方上下文文档流程控制文档用于辅助理解,访问日期为 2026-09-23;实现细节以固定 commit 为准。

前一篇:《Rill Flow 执行机制解析:DAG 调度、异步回调与失败重试》。下一篇将讨论状态保存、并发推进与业务幂等,把“定义能够表达什么”继续推进到“故障之后能够保证什么”。

按时间浏览