一个视频处理流程,为什么需要工作流引擎?
假设要实现一个视频处理接口:先探测媒体信息,再拆分片段,把各片段交给转码服务,最后合并产物。下面是贯穿本文的教学案例,不代表真实生产链路或个人项目经历。
探测 → 分片
├→ 转码 0 ─┐
├→ 转码 1 ─┼→ 合并
└→ 转码 2 ─┘
如果这些步骤很短、按顺序执行,一段普通业务代码就可以完成。复杂性来自执行过程中的变化:转码服务先返回任务编号,几分钟后才生成文件;第二个分片失败,其他分片已经成功;负责协调的服务在等待期间重启,内存里的进度随之消失。
这时,需要有人回答:哪些结果已经可用?失败节点还能尝试几次?完成通知重复到达,会不会再次发起合并?工作流引擎的价值,就在于集中管理依赖、状态、等待和后续推进,让这些规则有明确的表达与检查入口。
但引入引擎也意味着增加配置、存储和运维成本。选型之前,应先理解不同引擎如何组织执行,再判断其机制与业务是否匹配。本文比较公开文档和源码中的设计,不包含运行测试或性能基准。
同样叫 Workflow,几种引擎如何表达执行?
先选取三种有代表性的引擎作为参照。下面比较的是它们的主要抽象与设计侧重点,各项能力之间存在重叠。
| 引擎 | 流程表达 | 执行进度如何组织 | 设计侧重点 |
|---|---|---|---|
| Airflow | Python 定义任务依赖 | Dag Run 与 Task Instance | 批处理调度与运行管理 |
| Flowable | BPMN 流程模型 | 流程实例、网关与等待状态 | 业务流程及人机协作 |
| Temporal | 代码定义 Workflow | 事件历史与执行重放 | 持久化执行与故障恢复 |
| Rill Flow | 声明式 DAG | 节点状态、派发与完成处理 | 分布式服务编排与异步协作 |
Airflow 把流程定义成 Python 代码,调度器根据任务依赖和运行状态安排执行。对于周期性数据处理,读者往往首先关心某次运行对应哪些数据、失败任务如何重跑、历史区间如何补跑。它也支持手动及事件触发,不能简单等同于定时任务框架。
Flowable 用 BPMN 表达服务任务、人工任务、网关等业务元素。假如视频发布前需要人工复核,流程到达人工任务后进入等待状态,审核结果决定后续分支,这种表达很贴近业务规则。它同样能执行自动化服务调用,并不限于审批。
Temporal 则让开发者用代码表达流程,并通过事件历史与重放恢复执行进度。这里需要遵守 Workflow 代码的确定性约束,外部交互通常放到 Activity 中。其价值适用于需要持久化执行的流程,不取决于任务一定持续几小时或几天。
Rill Flow 将流程触发、流程编排和任务执行分开:定义描述依赖,引擎寻找可运行节点,再把任务交给执行器。它最初面向微博视频业务的复杂流程与并发任务,媒体处理是理解其设计的自然入口。官方简介
因此,Workflow 不等于 DAG。BPMN、程序代码与图模型可以表达不同的控制逻辑。比较时应关注执行语义,而不能只数有没有条件分支、重试按钮或可视化页面。
工作流引擎共同需要回答的五个问题
定义如何变成一次执行?
“先分片再转码”是流程定义;处理视频 A 是一次执行实例;视频 A 的第二个转码分片是具体任务实例;这个任务失败后再次调用执行器,又是一次尝试。它们需要不同的身份,排障时才能区分“两个视频同时处理”和“同一分片重复执行”。
Airflow 明确区分 Dag Run 与 Task Instance;Flowable 从流程定义创建流程实例;Temporal 从 Workflow 定义启动执行。它们的对象模型不同,但都需要把可复用的定义与本次运行的进度分开。对视频接口而言,只记录流程名称,无法定位某个失败分片。
下一步什么时候可以运行?
合并节点应等待全部必需的转码结果。这里“任务结束”和“任务成功”具有不同含义:失败也是一种结束,却不能据此生成完整视频。条件分支还会带来另一种情况:某条分支根本不应执行,汇合时需要遵守相应的分支语义。
Airflow 通过任务依赖与触发规则控制运行条件;Flowable 用网关表达路由与汇合;Temporal 可以在 Workflow 代码中等待所需结果。Rill Flow 则结合依赖和节点状态寻找可运行任务。阅读各自实现时,应先明确成功、失败、跳过如何影响后继,不能把所有箭头都解释为“前一步返回就调用下一步”。
结果如何传递?
“转码之后合并”描述控制关系,“合并读取哪些文件”描述数据关系。顺序正确不代表数据完整:如果多个分片都覆盖同一个结果字段,合并仍然可能只拿到最后一个文件。
Airflow 提供 XCom 传递任务间信息;Flowable 使用流程变量;Rill Flow 用输入输出映射连接节点与上下文。具体存储和作用域不同,都需要明确结果属于哪个实例、哪个分片。对于视频文件,合理的业务设计通常是传递对象地址、分片序号和校验信息,把大文件留在外部存储;这是应用设计建议,并非所有引擎的默认行为。
异步等待怎样结束?
转码接口返回任务编号,只能说明请求已受理。引擎还需要另一个信号,确认业务结果已经产生,才能推进合并。等待可以通过外部事件、轮询、人工完成或执行器通知表达,具体机制取决于框架与接入方式。
Flowable 的人工任务是明确的等待状态;Airflow 的 Sensor 用于等待外部条件;Rill Flow 的异步函数节点在派发之后等待完成处理。接口适配必须保留“受理”和“完成”的区别,否则流程图看起来正确,运行时却可能提前读取不存在的产物。
失败后依据什么继续?
第二个分片失败时,节点重试决定是否重新调用转码服务;协调进程重启后,状态恢复决定从哪里继续判断;Temporal 的重放则根据事件历史重新建立 Workflow 执行进度。它们解决的问题不同,保存过节点状态不等于已经具备事件重放能力。
还存在一个独立问题:转码其实已经成功,只是响应丢了。再次调用可能重复生成文件、计费或发送通知,因此业务操作仍需幂等。Temporal 也明确建议 Activity 设计为幂等。对本例,可以由执行器使用稳定的业务键识别同一视频、分片和处理参数,记录已有结果。引擎的重试策略不能替代这项业务约束。
Rill Flow 的四项设计选择
以下分析固定到公开仓库的 5cded0b 版本。这里只讨论普通非流式输入、未启用关键路径提前完成的执行方式;收益与代价是本文基于实现作出的归纳。
声明依赖与数据映射
Rill Flow 支持 YAML 和可视化编排。DAGStringParser 将描述解析为模型,依赖关系与数据映射随后参与执行。对视频案例,next 表达后继关系,inputMappings、outputMappings 说明节点如何读取输入、保存输出。
这种设计让流程结构集中可见:新增一个媒体检测步骤时,能够明确检查它处于哪两个节点之间、需要哪些数据。对已经独立部署的服务,流程调整与计算逻辑有了各自的维护入口。
代价是配置本身也成为需要管理的软件。YAML 语法正确,并不能证明服务地址可用、字段类型匹配或新旧输出兼容。配置发布前仍需检查图结构、数据契约和业务条件;已有运行实例如何关联定义版本,也应在接入时验证,不能假设修改配置会安全改变全部在途任务。
分离编排与具体计算
DAGOperations 按任务类别组织执行,函数节点由 FunctionTaskRunner 处理,任务再经派发路径交给执行器。官方架构列出 HTTP、Serverless Function 和扩展执行器等接入形式。
这对异构处理链路有实际意义:探测服务与转码服务可以使用不同语言,模型服务也可以拥有自己的资源环境。编排层负责决定哪些工作可执行,具体计算由对应服务承担。执行器能够独立演进,但“支持接入”不代表引擎自动替它分配最合适的 CPU 或 GPU。
接入成本集中到了边界协议:怎样表达受理结果、成功与失败,回调如何关联原任务,超时后怎样查询产物。统一这些约定,才能持续复用执行器;仅把服务地址填进流程定义,还不足以形成可靠协作。
围绕状态与完成事件推进
DAGTraversal 根据当前状态寻找可运行节点,并在加锁的推进路径中先保存选中任务的待执行状态,再交给执行路径。函数节点完成处理后会触发后续遍历。简化后的路径如下,状态与上下文贯穿各阶段:
流程定义 → 解析校验
↓
判断可运行节点 → 派发执行
↓
完成处理 → 保存结果与状态
↓
再次判断可运行节点
不同分片可以在不同时间结束,每次有效完成都成为重新判断依赖的机会。这种方式适合协调异步任务,也让“为什么还不能合并”可以落实到具体节点与结果。
同时,保存状态和调用外部服务不是一个跨系统原子动作。进程可能在两者之间退出,执行器也可能重复通知。不能仅从加锁和状态保存推导出自动补发、自动补偿或业务只执行一次。相应故障窗口、通知协议和恢复入口,需要结合具体实现与部署验证;详细处理路径留到第二篇。
显式表达并行与汇合
ForeachTaskRunner 根据输入集合组织多组子任务与子上下文。各分片的输出写入所属组,再由父节点汇总;官方并行异步示例 展示了这种结构,但示例处理数值,不是可直接使用的视频转码方案。
其收益是把“对每项执行”和“收集所有结果”变成明确的流程结构,减少业务代码中分散的计数、等待与结果拼装逻辑。对于要求所有分片成功的案例,某个分片失败就必须继续处理,不能通过跳过让缺失文件变成有效结果。
并行规模仍需治理。集合长度、同时执行数量、上下文体积都影响资源消耗;合并服务还应按序号排序并检查缺片,不能把 URL 列表等同于业务完整性证明。
| 设计选择 | 适用场景中的收益 | 工程代价 |
|---|---|---|
| 声明依赖与映射 | 流程结构与数据关系集中可见 | 配置、契约与版本管理 |
| 分离编排与计算 | 复用不同技术栈的已有服务 | 执行协议与错误语义统一 |
| 状态与完成事件推进 | 协调完成时间不同的任务 | 幂等、超时及恢复路径验证 |
| 显式并行与汇合 | 组织集合处理和结果收集 | 并发、上下文与完整性治理 |
据此,本文将其设计理念概括为:编排层管理依赖、状态和推进,执行器承担计算,再通过数据映射与完成通知连接两者。 这是对公开实现的归纳,不是项目作者的原话,也不是 Rill Flow 独有能力的清单。
什么情况下值得选择 Rill Flow?
如果团队已经有独立服务,需要经常调整它们的依赖,并处理异步等待、集合展开和结果汇合,Rill Flow 值得重点评估。它的优势在于这些设计与需求贴合,实际收益仍取决于接入成本和运行方式。
如果核心需求是数据批处理、历史补跑以及数据系统集成,先评估 Airflow 的运行管理与生态;如果人工任务、业务流程沟通及 BPMN 语义占主导,先评估 Flowable;如果希望通过代码表达复杂控制流,并依赖事件历史恢复执行,先评估 Temporal,同时理解其确定性约束与 Activity 边界。
对只有几个短小顺序调用的接口,普通业务代码可能已经足够。引擎增加的状态存储、部署和排障入口,应该换来实际需要的能力,而不能仅靠“流程看起来更复杂了”证明价值。
Rill Flow 官方将高并发、低延迟作为设计目标,但这不构成它比其他引擎更快的证据。真正选型时,可以用同一条业务链路检查任务规模、调度等待、重试行为和故障恢复,再衡量维护成本。只比较示例是否能运行,或者把不同环境下的宣传数字放在一起,都不足以支持性能结论。
从设计理解走向源码阅读
把视频分片换成文档,把转码换成模型调用,依赖、数据传递与异步等待仍然存在。变化在于还需定义输出质量、调用预算和循环退出条件。LangGraph 对有状态 Agent 编排、持久化与人工介入的关注,可以作为后续讨论的参照。
接下来沿着一次具体执行追问:定义怎样进入运行状态,回调怎样推进后继,失败又如何触发重试?系列下一篇《Rill Flow 执行机制解析:DAG 调度、异步回调与失败重试》将展开这些源码路径。
参考资料已随对应论点链接。官方文档访问日期为 2026-09-23;Airflow 页面标注版本为 3.3.2,其他文档采用当日页面;Rill Flow 源码与示例均固定到上述 commit。本文没有实际运行四种引擎或进行压力测试。