当前位置:首页 > 工具 > 正文内容

Flink Agents 核心架构解析

访客 工具 2026年9月5日 1

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));
        }
    }
}

确保故障恢复时的数据一致性。

返回列表

上一篇:智能合约开发实践:代币与DApp构建

没有最新的文章了...

相关文章

Trojan服务器搭建与配置

一、整体架构(先对齐认知)Clash Meta (PC / iOS / Android)        ↓ TLS   Trojan Server (443)        ↓     InternetTrojan 的核心是: TLS + HTTPS 流量伪装 看起来像正常网站 非常适合...

Tailscale 的详细用法

Tailscale 是一种基于 WireGuard 协议 的 零配置 VPN(虚拟私有网络)服务,让设备之间能够 安全、加密地直接连接,就像它们在同一个本地网络一样。它的核心特点是 简单、安全、跨平台。Tailscale 非常适合 没有公网 IP、两台电脑不在同一局域网 的场景。 简单来说,Tailscale 是什么?Tailscale 是一款让你的各种设备(电脑、服务器、手机...

Clash Tun 模式 导致 爱快(iKuai SD-Wan)内网域名无法访问

一、Clash  DNS 配置dns:  enable: true  listen: 0.0.0.0:53  ipv6: true  enhanced-mode: redir-host  nameserver:    - 223.5.5.5    - 223.6.6.6iKuai 内网域名 ...

深入解析Node.js运行环境与异步I/O架构

深入解析Node.js运行环境与异步I/O架构

核心定义与价值Node.js本质上是一个JavaScript运行环境,而非编程语言或应用框架。它赋予了JavaScript脱离浏览器在服务端、命令行工具及网络应用中执行的能力。其核心意义在于:用单一语言打通前后端开发壁垒。基于事件驱动与非阻塞I/O的架构特性,Node.js在处理API网关、实时通信及微服务等I/O密集型场景时表现卓越,已成为现代后端工程的主流选择。浏览器沙箱限制1995年Java...

ADO.NET SQL参数化查询的最佳实践

在 ADO.NET 中执行 SQL 查询时,参数化查询是一种关键的安全措施和性能优化手段。它通过将 SQL 命令和用户提供的数据分开处理,有效防止了 SQL 注入攻击,并有助于数据库缓存执行计划。下面总结了几种常用的参数化查询方式。 1. 使用 SqlParameter 对象(推荐) 这是最推荐的参数化查询方式。通过显式创建 SqlParameter 对象,您可以精确控制参数的类...

基于ELK的日志集中化分析系统搭建

构建统一日志管理平台的必要性 在分布式架构中,各服务节点独立运行,日志分散存储于不同主机。传统通过命令行工具如grep、awk逐个检索日志的方式,在数据量庞大时效率极低,难以实现快速定位问题。为提升运维效率,需建立集中式日志处理体系,具备日志采集、传输、存储、分析与告警能力。 ELK技术栈核心组件解析 Elasticsearch:分布式搜索引擎,支持全文检索、实时数据分析和高可用集群部署,...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。