这个系列写到第 6 篇了。前面我把 Agent 怎么拆计划、ActionScheduler 怎么排依赖、状态存储怎么回放都聊过一遍,今天终于轮到真正干活的 ActionTask。
说实在话,ActionTask 在整个 Flink Agents 体系里不是最亮眼的组件。它没有调度器统筹全局的威风,也没有状态存储兜底的安心感,但 Agent 到底能不能把活干完,最后都得落到 ActionTask 一个一个动作的执行上。说白了,它就是那个“真正干活的类”。
接下来我会把它拆到骨头里:核心字段、模板方法、和 ActionContext 的配合方式、超时重试怎么调,再加上我实际跑下来遇到的一堆坑。如果你已经在搭自己的 Agent 执行引擎,或者自己写了一堆类似 Task 的类但总觉得越写越乱,这篇应该能帮你把职责理清。
先说清楚,这套 Flink Agents 并不是 Apache Flink 内核的一部分,而是构建在 Flink 实时计算底座之上的一套 Agent 执行框架。别把它和 DataStream API 里的 Task 搞混了:Flink 的 Task 跑的是算子,ActionTask 跑的是业务流程里的动作。前者是流式计算的物理执行单元,后者是一个可以被调度、重试、记录和编排的业务动作,完全两个层次的东西。
1. ActionTask 到底解决什么问题
我最早接触这套框架时也有个疑问:Agent 已经有了,为什么还要单独抽象一个 ActionTask?直接写成 Agent 里的一个方法不行吗?还真不行。你设想一下,一个 Agent 要完成一个目标,往往要拆成十几个动作:先拉配置、再等依赖数据、再调下游系统、最后写结果回执。如果在 Agent 里把这些逻辑用 if-else 和 for 循环写死,代码会变成一团没人敢动的面条。
ActionTask 做的事情,就是把“执行一个动作”这件事里所有横切关注点收拢到一个抽象基类里。状态怎么流转、超时怎么判定、失败了怎么重试、重跑怎么保证幂等、结果怎么回传,这些逻辑不该让业务开发者每次重复手写,而是全部下沉到 ActionTask 这个公共载体里。业务方只需要继承它、实现一个 doExecute 方法,剩下的交给框架。
1.1 三层职责:Agent、Scheduler 和 ActionTask 各管一摊
要理解 ActionTask 的定位,最好把整个框架分层看。我一般用一张表给新同学讲这件事:
| 组件 | 职责 | 生活类比 |
|---|---|---|
| Agent | 接收目标、拆解计划、对外汇总结果 | 项目负责人 |
| ActionScheduler | 管理依赖关系、调度顺序、并发控制 | 施工调度员 |
| ActionTask | 执行单个动作,管理状态和重试 | 具体干活的师傅 |
| ActionContext | 动作之间共享结果、错误和指标 | 手递手交接单 |
这套分层的核心思想是“依赖倒置”。Agent 不需要知道某个动作具体怎么实现,它只认 ActionTask 的通用接口:你有一个 actionId、一个优先级、一个执行入口,那我就能给你排队、给你调度、给你记录结果。动作的具体细节藏在子类的 doExecute 里,Agent 和 Scheduler 完全不关心。
这样做换来的好处是:以后新增一个动作类型,比如“调用 HTTP 接口”“从 Hive 读数据”“给下游系统发消息”,你只需要新写一个继承 ActionTask 的子类,再往 ActionRegistry 里注册一个 actionType。既有的一整套调度、重试、状态存储、监控能力直接就复用了,不用改任何框架代码。这就是抽象出 ActionTask 的最大价值。
1.2 一个 ActionTask 的完整生命周期
一个新 ActionTask 从创建到结束,大致走六个步骤,每一步都有它存在的理由:
- 创建阶段:执行计划解析后,每个 ActionDescriptor 通过 ActionRegistry 找到对应的工厂,实例化 ActionTask 并注入 actionId 和 ActionContext。
- 入队阶段:Scheduler 检查依赖。所有依赖动作都变成 SUCCEEDED 状态,当前任务才允许进入 ready 队列。
- 调度阶段:线程池取到可用线程后,开始执行 run() 方法。
- 执行阶段:模板方法调用 doExecute(),真正干活的代码在这里运行。
- 回调阶段:执行成功后触发 onSuccess,失败触发 onFailure,无论结果如何最后触发 onComplete。
- 记录阶段:finally 里把状态、执行耗时、结果写入状态存储,供后续回放、监控和失败恢复使用。
看起来很简单,对吗?但真正把它落地成一个能扛住生产流量的执行单元,坑全在后面。尤其是状态流转怎么锁死、重试怎么判定、回调怎么设计,这三块最容易出问题。下面逐个拆。
2. 核心源码拆解:字段、模板方法和执行上下文
源码解读不能只看注释,得看这个类到底怎么设计的。我从工程里把去掉日志、监控和 SPI 扩展之后的骨架代码拎出来讲,类全名是 io.flink.agents.task.ActionTask。先看它管着哪些字段。
2.1 核心字段和执行状态,一眼看懂这个类在管什么
ActionTask 的字段并不复杂,但每一个都是经过生产验证的:
| 字段 | 类型 | 说明 |
|---|---|---|
| actionId | String | 动作唯一 ID,幂等判断和依赖判断都靠它 |
| actionType | String | 动作类型,ActionRegistry 据此定位对应工厂 |
| priority | int | 调度优先级,数字越小越先执行 |
| timeoutMs | long | 单次执行超时上限 |
| maxRetries | int | 最大重试次数 |
| retryCount | AtomicInteger | 当前重试计数,线程安全 |
| status | volatile ActionStatus | 当前状态,多线程可见 |
| context | ActionContext | 执行上下文引用 |
| callbacks | CallbackChain | 成功、失败、完成三组回调 |
这里最值得说的是状态枚举。ActionStatus 的状态流转被刻意限制住了:PENDING 是初始态,RUNNING 表示正在执行,重试时会暂时回到 RETRYING,最终要么 SUCCEEDED,要么 FAILED,还有一个 CANCELLED 由外部取消触发。你会发现没有任何一条路径能从 RUNNING 跳回 PENDING,一个动作不可能回头重新排队。这个限制是有意为之。如果允许随便回退,调度器的依赖判断就乱了:下游都以为上游还在排队,整个 DAG 永远走不完。
2.2 为什么 run() 必须加 final,子类只能写 doExecute
下面这段是我去掉监控和日志之后的简化骨架,重点看 run() 方法的设计:
public abstract class ActionTask<T> implements Runnable { protected final String actionId; protected final ActionContext context; private final ActionOptions options; protected volatile ActionStatus status = ActionStatus.PENDING; @Override public final void run() { if (!status.compareAndSet(ActionStatus.PENDING, ActionStatus.RUNNING)) { return; } long start = System.currentTimeMillis(); try { T result = doExecute(); status = ActionStatus.SUCCEEDED; context.putResult(actionId, result); onSuccess(result); } catch (RetryableException e) { scheduleRetry(e); } catch (Throwable t) { status = ActionStatus.FAILED; context.putError(actionId, t); onFailure(t); } finally { context.record(actionId, status, System.currentTimeMillis() - start); onComplete(); } } protected abstract T doExecute() throws Exception; protected void onSuccess(T result) {} protected void onFailure(Throwable error) {} protected void onComplete() {} }run() 加 final 是我在这套框架里最喜欢的设计。一旦 run() 可以被重写,100 个动作里至少有 30 个会忘了更新状态、忘了调用回调,最后状态机一塌糊涂。现在基类把流程锁死,子类唯一必须实现的就是 doExecute,最多再覆盖 onSuccess、onFailure、onComplete 三个钩子。这就是模板方法模式:框架定好舞步,你填具体动作。
这个模式生活里也到处都是。公司里的报销流程就是一个典型:主管审批、财务复核、打款,这些步骤是制度定死的,你只需要填费用明细这一张表。如果让每个员工自己设计报销流程,公司早乱了。ActionTask 的 run() 就是那个被制度锁死的流程,doExecute 就是你填的那张表。如果你是刚接触设计模式的新手,不用记名词,理解成“基类把执行流程锁死,子类只填具体动作”就够了。
2.3 ActionContext:动作之间靠它传结果,别再用静态变量
ActionContext 大概是整个框架里最不起眼却最容易用错的东西。它的作用说白了就三条:传结果、传错误、传指标。我见过有人在多个 ActionTask 之间用 static 变量传数据,排查问题的时候被各种隐藏状态折磨到怀疑人生。用 ActionContext 显式传递,谁写了谁读了,一眼就能看清数据流。
public class ActionContext { private final Map<String, Object> results = new ConcurrentHashMap<>(); private final Map<String, Throwable> errors = new ConcurrentHashMap<>(); private final MetricRegistry metrics = new MetricRegistry(); private final StateStore stateStore; public void putResult(String actionId, Object result) { ... } public <T> T getResult(String actionId) { ... } public void putError(String actionId, Throwable error) { ... } public void record(String actionId, ActionStatus status, long costMs) { ... } }注意这里的 Map 必须用 ConcurrentHashMap。ActionTask 是在线程池里并行执行的,同一个上下文会被多个线程同时读写,普通 HashMap 在并发写入时轻则丢数据,重则直接死循环卡死 CPU。这一点没有商量余地,凡是多线程共享的容器,一律先考虑并发安全。
3. 完整实操:从构建 ActionTask 到调度执行
光看定义不够,我们直接上手,看一个 ActionTask 从构建、入队到执行完毕的完整过程。
3.1 用 Builder 构建一个能落地的 ActionTask
实际编码时我们不会手写 new ActionTask,而是通过 Builder 来构造。一个 HTTP 调用动作的构建长这样:
ActionTask<HttpResp> task = new HttpCallActionTask.Builder() .actionId("notify_order_center") .actionType("http") .url("https://oms.example.com/notify") .method("POST") .payload(payload) .priority(10) .timeout(Duration.ofSeconds(30)) .maxRetries(3) .baseInterval(Duration.ofSeconds(2)) .dependentOn("wait_order_center_ready") .build();为什么用 Builder?因为这个类字段太多,十来个字段如果全塞构造方法,参数顺序稍微错一位就是线上事故。Builder 的另一个好处是可以在 build() 里做必填校验:比如 actionId 为空直接抛出 IllegalArgumentException,把错误挡在启动阶段,而不是等运行期跑到一半才炸。这对排错体验非常友好,所有“这个对象不该被造出来”的情况,都应该尽早暴露。
3.2 线程池怎么配,才不会互相拖垮
很多人第一次接入这套框架会问:Flink 已经用自己的线程池跑算子了,ActionTask 还需要单独搞线程池吗?需要。Flink 的算子通道负责的是流数据的传输和处理,ActionTask 是 Agent 引擎里面的动作执行队列,两者任务来源不同、并发模型不同、生命周期也不同。ActionTask 的线程池是和 ActionScheduler 配套的,核心参数大致长这样:
ThreadPoolExecutor executor = new ThreadPoolExecutor( corePoolSize, maximumPoolSize, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(queueCapacity), new ThreadFactoryBuilder().setNameFormat("action-task-%d").build(), new ThreadPoolExecutor.CallerRunsPolicy() );两个关键点一定要记住。第一,队列不要用无界队列。无界队列唯一的作用就是内存悄悄耗尽,然后整台机器和你一起痛苦。给个固定容量,满了就触发拒绝策略,让上游慢下来,这才健康。你可以把它理解成 Flink 的背压:餐厅做不过来,门口排点队没问题,但别把大厅设计得无边无际。
第二,拒绝策略强烈建议用 CallerRunsPolicy。这种策略下,队列满了不会丢任务,而是由提交任务的调度线程自己来执行。牺牲一点调度效率,换取“任务绝不静默丢弃”。我在好几次下游抖动的时候,全靠这个策略保住了任务不丢失。
3.3 成功、失败和重试的判定逻辑
2.2 的骨架里出现了 scheduleRetry,这里单独展开。它的实现大致是:
private void scheduleRetry(Throwable e) { if (retryCount.get() < options.maxRetries()) { retryCount.incrementAndGet(); status = ActionStatus.RETRYING; long waitMs = Backoff.exponential(options.baseIntervalMs(), retryCount.get()); scheduler.schedule(this, waitMs); context.metrics.counter("ActionTask.retry").inc(); } else { status = ActionStatus.FAILED; context.putError(actionId, e); onFailure(e); } }这里最关键的设计是区分可重试异常和不可重试异常。框架里专门做了一个 RetryableException,只有抛出这个异常,基类才允许重试;其它 Throwable 一律直接失败并触发 onFailure。为什么要区分?因为盲目全量重试是给自己埋雷。
比如 Flink 写 Hive 表偶发的元数据锁冲突,这种属于典型可重试异常,重试两次加个短退避就能过去;但如果表 schema 根本不匹配,重试一万次结果都一样,纯属浪费资源。再比如参数校验失败的请求,重试也一样会失败。把“能不能重试”的判断交给异常类型,而不是 catch Exception 一把梭,这是生产级代码的基本修养。
3.4 回调链:结果通知下游的正确姿势
在 2.2 的骨架里,onSuccess 并不是直接通知外部系统,它只是留了个钩子。实际业务中,回调链长这样:
task.setCallbacks( result -> { context.mark("order_center_notify_done"); pushToMq(result); }, error -> alertService.send(error) );回调里有三个坑,今天我一起说掉。
第一,回调里不要做重试。重试是基类的事,回调只负责通知和标记,职责才能清晰。有人在 onFailure 里又发一次失败请求,结果和基类的重试机制叠加,下游收到好几条重复消息,很难查。
第二,回调里不要做耗时操作。回调是执行线程的一部分,一个回调阻塞 2 秒,线程池吞吐就肉眼可见地下降。如果回调要写库、发消息,丢进一个异步队列慢慢处理。
第三,onComplete 里的清理动作要小心。比如释放连接、关闭句柄,一定要放在 try-finally 里。否则前面某个清理动作抛了异常,后面的清理代码全被跳过,资源泄漏就是这么发生的。
4. 参数计算与调优:像调 Flink 作业一样调 ActionTask
读到这里,一个 ActionTask 的基本运转方式已经清楚了。下面这部分可能是大家最关心的:参数到底怎么定?我见过太多项目,超时拍脑袋写个 10 秒,重试写 3 次,结果线上动不动报警,也不知道该调什么。
4.1 超时、重试次数和退避间隔怎么算(附真实数值例子)
参数不是拍出来的,是压测数据喂出来的。我每次接入新动作,都会先要求一份下游接口的压测数据,然后套下面的估算逻辑。
拿一个 HTTP 通知下游的场景举例。压测数据通常是这样的:TP50=120ms,TP99=800ms,最坏的 P99.9=2.1s。
- 超时定多少?我一般用 TP99 × 5,原因是要给偶发抖动留缓冲。800ms × 5 = 4s,所以 timeout 定 4 秒,而不是图上看到的 2.1 秒。太贴 P99.9,一次 GC 抖动就会误杀正常请求。
- 基础重试间隔 = 平均耗时 × 2。平均约 200ms,乘 2 得到 400ms,工程上取整 500ms,作为 baseInterval。
- 重试等待序列按指数退避:第 1 次等 500ms,第 2 次等 1s,第 3 次等 2s,三次重试的总等待时间是 3.5s。
- 最坏情况算一笔账:每次执行都卡到 4s 超时,三次执行 12s,加上等待时间 3.5s,一共 15.5s。
- 如果业务 SLA 是 30s,那没问题;如果 SLA 只有 10s,就把 maxRetries 压到 2 次,或者把超时降到 2s 再算一遍。
这套计算方式看起来朴素,但它把每个参数都挂到了数字上。以后线上出问题,你至少知道先看哪个指标、改哪个参数会影响多少耗时。比拍脑袋强太多。
4.2 资源隔离与并行度:别把所有动作塞进一个大池子
刚开始接这套框架时,我把所有 ActionTask 塞进同一个线程池,结果写库动作把 HTTP 调用池子挤爆了,系统表现就是调用下游大面积超时。后来我改成按动作类型分池:http-pool、jdbc-pool、hive-pool 各一摊。虽然线程总数多了一点,但故障面被切小了。某个池子被打爆,最多影响这一类动作,不会全局雪崩。
池子大小怎么定?我的经验是分两类看:
- 计算型动作(解析 JSON、做聚合、跑规则引擎):核心线程数不超过 CPU 核数 × 2,线程再多也没用,CPU 已经饱和。
- 阻塞型 IO 动作(调 HTTP、读数据库、写 Hive):核心线程数可以到 CPU 核数 × 4 甚至更高,因为线程大部分时间在等待网络返回,占用的计算资源很少。
还有一个细节:同一个 actionId 同一时刻只能有一个线程在执行,这个靠 2.2 里的 CAS 判断保证。如果你自己写框架,务必保留这个判断。否则并行度一高,同一个任务被重复消费,幂等性再强也顶不住。
4.3 幂等设计是 ActionTask 的保命底线
这套框架每天可能要重跑无数个动作,加上宕机恢复、手动补偿,一个动作被重复执行是非常正常的事。没有幂等保护,重试机制越完善,事故越容易从重试里冒出来。幂等键的设计一般长这样:
String dedupKey = actionId + ":" + executionId; if (stateStore.exists(dedupKey)) { log.info("action already completed, skip. key={}", dedupKey); return; } T result = doExecute(); stateStore.markCompleted(dedupKey, result);核心逻辑就三句:先查记录,发现已成功就跳过;没成功才执行;执行完立刻标记。这里要特别提醒,幂等键不建议只用 actionId。同一个动作在不同轮次里可能会带不同参数,只用 actionId 会误伤。actionId + executionId 的组合才能精确表达“这一轮执行里的这个动作”。
很多团队栽的跟头在于只做了执行器内部的幂等,忘了外部系统也要幂等。比如你调一个下游接口,请求里必须带上 dedupKey,下游才能做去重。否则你这边查到了已执行,下游那边可能已经被重复调用两次了。所以接 ActionTask 时,先问自己一句:这个动作重跑一次会不会产生不同结果?会,就必须想清楚幂等键怎么传、在哪一步标记完成。
5. 常见问题与排查技巧实录
这部分是实践中积累的排查经验,整理成速查表,先对症状再动手。
5.1 一张故障速查表,先对症状再动手
| 症状 | 可能原因 | 排查与处理 |
|---|---|---|
| 任务一直卡 PENDING | 依赖动作没成功,或自己没被调度 | 查状态存储里父动作状态,修复父动作失败原因后重放 |
| 重试后同一动作被重复执行 | 没有幂等键保护 | 加 dedupKey,执行前查 stateStore |
| 线程池队列持续堆积 | 下游动作太慢 | 看下游 TP99,调超时或单独分池 |
| 状态 SUCCEEDED 但业务没生效 | doExecute 没校验业务码 | 检查 HTTP 200 之外响应体的业务 code |
| 回调一直没触发 | 回调异常被吞 | 查 callbackError 日志,回调内 try-catch |
| 状态恢复后结果错乱 | 状态存储序列化问题 | 检查时间字段类型,统一序列化方案 |
这些问题的共同点是:现象在 ActionTask 层,根因往往藏在依赖、序列化或下游系统里。所以排查时不要只盯着 ActionTask 的日志,要顺着状态存储的记录往前追。状态存储里记了什么、没记什么,往往就是定位问题的第一现场。
5.2 我自己踩过的两个坑,希望你别再踩
第一个坑是重试间隔写死。最早我把重试间隔写成固定值 5 秒。某天下游服务故障恢复之后,我这边还在按固定节奏重试,白等了好几轮。后来改成指数退避加随机抖动,同时让重试前先探测下游健康检查接口,一旦恢复立刻重试,效果明显好了很多。指数退避能避免下游刚恢复时瞬间被打爆,随机抖动能避免多个任务同时重试造成一波流量高峰。
第二个坑是只校验 HTTP 状态码。当时我把“调用成功”理解成 doExecute 不抛异常,结果下游接口因为参数问题返回了 HTTP 200 但业务 code 是 500,ActionTask 状态一路 SUCCEEDED,数据却完全没生效。排查了很久才意识到,doExecute 里必须把响应体里的业务状态码也校验掉,业务不成功直接抛 RetryableException,而不是让异常被吞掉、静默成功。那次之后我对所有网络类动作的校验都多了一行:先看 HTTP status,再看 body 里的业务 code,两层都过了才算成功。
在上面那套框架上我大概跑了快一年的数据业务编排,最直观的感受不是用了多高级的技术,而是失败处理终于有地方去了。以前我们写批处理脚本,每段 shell 里都得自己处理异常、自己记日志、自己保证重跑安全,动作一多就是灾难。现在每个动作都是一个 ActionTask,状态、重试、回执全部内建,业务代码只需要关心 doExecute 里那一小段真正的业务逻辑。
如果你也要接入类似的 Agent 框架,我给个最朴素的小技巧:每写一个新的 ActionTask,先对着自己问三个问题——这个动作重跑一次会怎样?超时了怎么处理?失败回调通知谁?三个问题想明白了,代码基本不会出大事故。我每次新接一个动作需求都这么干,你可以试试。