这是系列第五篇。第三篇解释数据与并行结构,第四篇讨论故障窗口。本篇把同一条教学视频链路接到已有转码服务,分析执行协议与运行治理。依据仍为公开仓库 5cded0b 版本;本文阅读了实现和测试源码,校验了示例语法,没有启动引擎或执行器,也没有进行负载测试。
已有业务服务怎样进入工作流
假设团队已经有一个转码服务:提交源文件后返回作业编号,后台完成计算,另有接口查询进度。引入工作流不必把视频处理逻辑重写进引擎。真正需要增加的是一层适配,把引擎的节点调用转换为业务作业,并把业务完成转换为引擎能够理解的通知。
这条链路至少包含四类职责。引擎判断哪一个任务现在可以运行;派发器选择协议并组织请求;执行器解释业务参数、调用转码能力并形成结果;结果存储保存媒体产物。控制状态、计算过程和大文件各有归属,不能把“任务已成功”当作媒体文件永远可用的保证。
已有服务若只支持查询,适配层还需要负责轮询和完成通知;若原生支持回调,则需要转换状态与结果格式。两种方式都必须明确作业编号如何关联到工作流任务。这里描述的是接入设计,并不意味着引擎会自动识别任意第三方作业协议。
DispatcherExtension提供资源处理扩展入口,FunctionTaskDispatcher根据资源协议选择具体处理实现。由此可以复用已有服务,而不是要求业务计算和引擎运行在同一进程。不过,协议可扩展只解决“怎样发起调用”,业务幂等、结果有效期和通知确认仍需要双方约定。
对于教学视频链路,探测和分片可以是短请求,转码通常更适合异步执行,合并是否异步取决于实际耗时。应该按任务特性选择模式,不必为了统一配置,把所有计算都塞进一个长期占用连接的同步调用。
一次派发包含哪些信息
resourceName 主要回答去哪里、按什么资源协议执行,不等于业务唯一键。函数派发器先准备请求头,再把配置中的 parameters 作为输入缺省值补入;已有同名输入不会被这一步无条件覆盖。随后根据资源描述选择处理器。配置评审应同时检查地址、协议与最终输入,不能只看节点名字。
输入映射发生在调用准备过程中,决定业务服务真正得到哪些字段。HttpInvokeHelperImpl又会把输入组织为请求体、查询参数和请求头,并在相应请求构建路径中加入执行实例等关联信息。不同资源处理器未必使用完全相同的请求形式,因此接入前应沿实际协议核对一次。
官方 Java 示例的 ExecutorController读取 X-Mode 和 X-Callback-Url,把请求体放入执行上下文。模式默认是同步;回调地址用于异步完成后的通知。下面只展示业务请求体的教学片段,不是完整 HTTP 报文:
{
"source_url": "segment-001.mp4",
"segment_index": 1,
"profile_version": "h264-v2",
"business_key": "video-A:segment-1:h264-v2"
}
前三个字段表达计算输入,最后一个是建议由业务增加的稳定键,不是该控制器自动生成的协议字段。执行实例定位一次流程,任务标识定位图中的节点,业务键定位需要去重的操作;一次调用尝试还可能需要自己的编号。这些标识有关联,但不能互相替换。
请求与结果也应约定版本。执行器升级后若把原来的结果地址改成对象,旧输出映射可能立即失效。可以由适配层维持原格式,也可以发布新的契约版本并同步修改定义。无论选择哪一种,都应保存可定位的版本信息,避免排障时只能猜测当时使用的是哪套参数。
回调地址属于协议的一部分,不能只在开发环境验证一次。执行器部署位置变化后,通知方向的网络连通性可能与请求方向不同。接入检查需要分别确认两条路径,并说明谁负责回调鉴权和错误响应的解释;这与视频算法是否运行正常是两个问题。
把受理与完成写进执行协议
ExecutorWrapper给出了很直观的对照:同步模式直接执行函数并返回结果;异步模式把工作提交给线程池,立即返回 result_type: SUCCESS。在异步请求的这个位置,成功表示受理路径完成,不能理解为视频已经转码完成。
后续异步计算正常结束时,包装器在结果中加入成功类型;异常时构造失败类型与错误说明,最后通过 TaskFinishCallback发送通知。它序列化的是结果对象本身。不要把 Java 包装对象的字段层级额外套到 HTTP 通知体上,也不要把流程里恰好名叫 callback 的业务节点当作这个完成接口。
引擎 Java 执行器 业务计算
│── 请求 ─────>│ │
│<─ 已受理 ───│── 提交计算 ───────>│
│ │<─ 结果或异常 ─────│
│<─ 完成通知 ─│ │
│── 保存结果并推进后继 │
以人工构造的转码结果为例,通知体可以表达如下内容。业务字段由适配层提供,真实地址与回调鉴权信息均省略:
{
"result_type": "SUCCESS",
"segment_index": 1,
"result_url": "transcoded-001.mp4"
}
失败通知则应清楚表达任务失败,而不是仍返回成功类型、只在某个描述字段里写“失败”。官方包装器使用 FAILED 和 error_msg 表达异常;业务还可以约定可分类的错误信息,但引擎是否据此重试,要继续看已配置的失败策略,不能由执行器自行假定。
这份示例适合作为协议入口,并不是完整的生产执行器。异步工作放入进程内线程池,进程退出后的作业保障需要另行设计;通知发送代码捕获异常后记录日志,所示路径没有持久化通知队列。其发送函数读取响应体,未在该处明确检查 HTTP 成功状态,因此不能把一条发送成功日志等同于引擎已确认完成。
业务接入可补充作业持久化、结果查询和通知重试,但应先定义通知的确认条件:哪些响应表示已接受,哪些表示重复完成,哪些需要稍后重试。重试应有退避与终止条件,并保留原结果。对已完成的计算重新发送通知,通常不应再次执行转码。
当执行器变慢或拒绝请求时
速率与并发解决不同问题。每秒只接受少量请求,不代表同时运行的作业很少;任务耗时变长,在途数量仍会持续增加。foreach 的分组并发控制又限定在相应父结构和配置条件内,不能直接充当所有工作流共享的转码资源配额。
固定版本的 DAGSubmitChecker会按服务或业务配置调用 TrafficRateLimiter。后者使用带秒级时间片的键和 Redis 脚本检查许可。这条调用路径控制流程提交流量,不是在每个转码线程执行前自动扣减全局并发额度。类名中的“限流”必须连同作用位置一起解释。
异常策略也影响保证:该限流器捕获异常后返回允许。这样可以避免检查组件故障直接阻断全部请求,但意味着它不是任何故障下都严格生效的容量硬边界。外部执行器仍需根据自身队列和资源情况做准入控制,不能只依赖入口流量配置。
DAGResourceStatistic包含资源状态和冷却信息处理,会从相应响应中的 sys_info 或 error_detail 读取 retry_interval_seconds。这是可供资源治理使用的信息,并非任意 HTTP 拒绝都会自动产生完整背压。实际效果还取决于使用的派发路径、检查开关和资源配置。
如果执行器队列已满,应明确拒绝是在受理之前还是之后发生。受理前拒绝可以稍后重试;受理后又返回模糊超时,可能留下一个已经运行的作业。业务应该查询已有作业或使用稳定键重复提交,避免把“忙”直接变成另一份重复计算。
背压还需要反馈回到生产者:限制新流程进入、限制展开速度,或延后派发。只增加线程数和重试次数可能放大拥塞。接入方案应为队列设置容量、为等待设置期限,并说明拒绝后的责任归属。这些是运行设计建议,本文没有用负载测试证明某个容量数值。
容量评估还要考虑多个流程共享资源。单个视频只有三片,多个视频同时到达时仍可能形成大量请求;限制每个父节点只运行一组,也不意味着整个服务只有一项计算。接入层需要选择配额的归属,例如按资源池、租户或业务等级分配,并明确闲置额度能否共享。这类全局策略应与局部并行配置分别记录。
怎样定位流程卡在哪里
排障可以沿着四个问题走:任务是否满足依赖,是否已经派发,执行器是否产生结果,完成通知是否被引擎接受。同样显示运行中的节点,可能分别在排队、计算、等待通知,或处理通知失败。只看一个总耗时指标,往往无法区分这些位置。
源码中已有请求日志和 Trace 接入。ContextTraceHook在跟踪开启时把跟踪标识写入相应上下文,OpenTelemetryTracer提供跟踪能力。是否启用、标识是否穿透业务服务、采集端是否保留数据,都需要部署时确认;存在接口不代表默认获得一条完整链路。
下面按排障目的整理观测信息,第三列明确区分已有入口与接入建议:
| 想回答的问题 | 需要关联的信息 | 依据或补充位置 |
|---|---|---|
| 哪次流程的哪个节点 | 执行实例、任务名 | 派发及运行日志已有相关入口 |
| 是哪一次远端调用 | 尝试编号、业务键、作业号 | 建议适配层建立关联记录 |
| 时间花在哪里 | 等待、派发、执行、通知延迟 | 建议分阶段记录时间戳 |
| 是否持续过载 | 拒绝数、在途数、队列深度 | 建议执行器与入口分别采集 |
| 哪份结果被采用 | 输入版本、产物版本、通知结果 | 建议结果存储保留提交记录 |
业务键适合查一项操作,却不适合直接成为高基数监控指标的标签。可以在日志或明细记录里保存它,在聚合指标里按资源、错误类别和流程版本统计。这样既能观察整体趋势,也能从异常实例继续追查具体作业。
日志内容也应有取舍。派发路径可能记录请求参数,而视频地址、签名信息和文档正文未必适合长期完整留存。接入时应确认脱敏与保留策略,并保存足够的摘要或引用用于排障。可观测性的目标是解释执行过程,不是无限复制业务数据。
上线前应验证哪些接入场景
最小验收不应只有一次正常调用。同步成功要验证返回结构与输出映射;异步成功要验证受理之后仍会等待,并在完成通知被接受后推进;业务失败要验证错误分类与配置策略一致。对于每一种场景,都应同时检查引擎状态和业务产物。
重复通知应验证不会造成额外的业务副作用,并确认发送端如何停止重发;慢执行应验证超时后原作业是否仍存在,迟到结果如何处理;过载应验证拒绝发生在受理之前,以及恢复容量后是否出现集中重试。每个预期都要有明确观察点,不能只看接口返回码。
验收记录最好保留一条可复查的关联链:哪份输入触发了哪个任务、派发到哪个业务作业、产物由哪个通知提交、后继何时开始。若某一步只能通过人工猜测关联,上线之后面对并发和重试会更加困难。补全这条记录比增加一张笼统的成功率图更能解释故障,也能帮助确认结果是否被重复使用。
还需要安排两种中断:执行器受理后退出,和结果已生成但通知方向断开。前者检验作业的保存与查询能力,后者检验通知恢复能力。示例线程池和单次通知代码不足以证明这两种恢复已经成立,需要用真实接入方案补足并验证。
阅读 Java 包装器测试源码可以看到同步调用与异步提交的断言,它帮助确认示例的分流意图,但不能代替上述故障验证。本文只阅读这些断言,未实际运行测试。运行指标、容量与告警阈值也应根据接入环境测量,不在文章中填写未经验证的经验数字。
官方执行器文档和过载保护文档作为辅助资料,访问日期为 2026-09-23;实现解释以固定版本源码为准。前一篇是状态保存、并发推进与业务幂等。下一篇将把任务从视频转码换成模型调用,继续讨论结构正确之外的质量与预算问题。