Flink Agents 核心架构解析
0x00 架构概述
Flink Agents 框架的核心设计是"事件驱动 + 状态隔离 + 多语言协作":通过 Agent 和 AgentPlan 实现业务逻辑的声明式定义,结合 Flink 的分布式与高并发能力,确保任务的高效执行。同时支持 Python 工具和模型的集成,兼顾开发灵活性与运行效率,适用于复杂 AI 代理任务的分布式执行。
该框架并非重新发明轮子,而是对 Flink 原生组件在 Agent 场景下的语义化封装。本文将介绍其核心组件,对比其与 Flink 原生组件的映射关系,并通过实例说明其设计精髓。
0x01 核心组件结构
可以将 Flink Agents 的执行流程类比为"制作一道菜":
1.1 核心组件
Flink Agents 基于 Flink 流处理能力构建,包含四个主要组件:
- Agent(顶层设计,定义"做什么"):用户定义的智能实体,包含动作(Action)和资源定义,明确业务逻辑。
- AgentPlan(中间编译层,确定"怎么做"):将 Agent 编译为可执行计划,定义动作触发规则和资源映射。
- ActionExecutionOperator(运行时执行层,负责"协调调度"):Flink 集群中的执行核心,负责接收数据、调度任务、管理状态。
- ActionTask(最小执行单元,负责"具体实施"):具体执行任务,分为 JavaActionTask 和 PythonActionTask。
这种设计使得系统具备良好的扩展性和可维护性。
1.2 组件成员映射
1.2.1 Agent 到 AgentPlan
| Agent 成员 | AgentPlan 对应成员 | 说明 |
|---|---|---|
| _actions(装饰器定义) | actions, actions_by_event | 通过 @action 定义的动作被编译到 AgentPlan 的动作映射中 |
| _actions(add_action 添加) | actions, actions_by_event | 通过 add_action 添加的动作同样编译到动作映射中 |
| _resources | resource_providers | Agent 中注册的资源被转换为资源提供者 |
1.2.2 AgentPlan 到 ActionExecutionOperator
| AgentPlan 成员 | ActionExecutionOperator 对应成员 | 说明 |
|---|---|---|
| actions | getActionsTriggeredBy() | Operator 根据事件类型查找对应动作 |
| resource_providers | RunnerContextImpl | 提供运行时所需的资源 |
| config | metricGroup, builtInMetrics 等 | 用于配置指标和其他运行时行为 |
1.3 执行流程
具体执行流程如下:
Action Code → Agent → AgentPlan → ActionExecutionOperator → ActionTask → Flink Runtime
以 ReActAgent 为例:
- 用户定义 ReActAgent,包含 start_action 和 stop_action
- 通过 AgentPlan.from_agent() 编译成计划
- AgentPlan 被传递给 ActionExecutionOperatorFactory 创建执行器
- ActionExecutionOperator 接收 InputEvent,依次执行 start_action → ChatRequestEvent → 内置动作 → stop_action → OutputEvent
0x02 与原生 Flink 的映射
Flink Agents 是对 Flink 流处理能力的领域封装,其核心组件均可映射到 Flink 原生组件。
2.1 映射关系表
| Flink Agents 组件 | 原生 Flink 对应组件 | 核心角色 |
|---|---|---|
| Agent | StreamGraph / 用户 DataStream 代码 | 高层业务逻辑声明(做什么) |
| AgentPlan | JobGraph | 编译后的可执行计划(怎么拆) |
| ActionExecutionOperator | KeyedProcessOperator / StreamOperator | 运行时核心执行算子(核心载体) |
| ActionTask | 算子内处理单元 / AsyncFunction 任务 | 原子执行任务(最小执行单元) |
2.2 组件对比
- Agent:用户定义的智能实体,与 Flink 用户编写的 DataStream 代码类似,定义"要做什么"。
- AgentPlan:将 Agent 编译为可执行计划,与 Flink JobGraph 类似,定义"如何执行"。
- ActionExecutionOperator:Flink 集群中的执行核心,与 KeyedProcessOperator 类似,负责运行时调度。
- ActionTask:最小执行单元,与 Flink 的 AsyncFunction 任务类似,支持异步执行。
0x03 实例解析
我们以"餐厅自动化服务系统"为例,说明 Flink Agents 的工作原理。
3.1 Agent:餐厅菜单和规则手册
- 定义餐厅能提供的菜品和服务(动作)
- 规定所需设备和食材(资源)
- 遇到什么情况执行什么操作(事件监听)
- 描述完整服务流程(start_action、stop_action)
3.2 AgentPlan:操作流程图
- 将菜单分解为具体步骤
- 明确触发条件(Event → Action)
- 准备资源清单(Resource)
3.3 ActionExecutionOperator:执行管理层
- 接收顾客订单(InputEvent)
- 根据流程图分配任务给员工(Action)
- 协调各岗位工作(管理状态)
- 确保流程顺畅(处理并发和容错)
3.4 ActionTask:具体服务步骤
- 服务员接到任务(ActionTask)
- 可能一步完成,也可能异步执行
- 完成后报告结果,触发下一步
3.5 示例流程
- 顾客进店(InputEvent)→ 接待员(start_action)
- 生成厨房订单(ChatRequestEvent)→ 厨师制作
- 完成通知(ChatResponseEvent)→ 收银员(stop_action)
- 顾客结账(OutputEvent)
0x04 并发与并行机制
Flink Agents 充分利用 Flink 的并发模型,确保高效执行。
4.1 Flink 并发模型
- 使用 KeyedStream 按 key 分区,确保相同 key 由同一实例处理
- ActionExecutionOperator 可在多个实例上并行运行
// CompileUtils.java
public static <IN, K> DataStream<Object> connectToAgent(
DataStream<IN> inputStream, KeySelector<IN, K> keySelector, AgentPlan agentPlan) {
return connectToAgent(inputStream.keyBy(keySelector), agentPlan);
}
4.2 Key 状态隔离
// ActionExecutionOperator.java
private transient ListState<ActionTask> actionTasksKState;
private transient ListState<Event> pendingInputEventsKState;
每个 key 维护独立状态,防止竞争。
4.3 邮箱线程模型
private final transient MailboxExecutor mailboxExecutor;
mailboxExecutor.submit(() -> tryProcessActionTaskForKey(key), "process action task");
实现协作式多任务处理,避免阻塞。
4.4 异步任务处理
public class ActionTaskResult {
private final boolean finished;
private final List<Event> outputEvents;
private final Optional<ActionTask> generatedActionTaskOpt;
}
支持异步延续任务。
4.5 检查点与容错
@Override
public void snapshotState(StateSnapshotContext context) throws Exception {
if (actionStateStore != null) {
Object recoveryMarker = actionStateStore.getRecoveryMarker();
if (recoveryMarker != null) {
recoveryMarkerOpState.update(List.of(recoveryMarker));
}
}
}
确保故障恢复时的数据一致性。
