Realtime ComputeFlink支援基於開源Apache Flink Agents架構開發事件驅動的流式 AI Agent 作業。本文介紹 Flink Agents 的核心概念、版本要求、快速上手路徑、自訂開發流程,以及作業提交與依賴管理。
概述
Apache Flink Agents是 Apache Flink 社區推出的全新子專案,提供了一個用於構建事件驅動型 AI Agent 的開發架構。它基於 Flink 久經考驗的流式引擎,將分布式、有狀態、容錯、串流的能力帶入 Agentic AI 領域,讓智能體真正走進 ToB 生產環境。
Flink Agents內建大模型調用、工具調用、記憶管理、動態編排、可觀測性等核心模組,開發人員可快速構建大規模、持續運行、端到端可靠的生產級 AI Agent。
Flink Agents具備如下核心優勢:
|
能力 |
說明 |
|
分布式協同 |
不是把單機 Agent 跑成多副本,而是從架構層提供事件路由分區、狀態分布一致、故障斷點恢複等能力。 |
|
有狀態記憶 |
基於 Flink State 託管 Agent 記憶,隨 Checkpoint 自動持久化、恢複與重分布,無需額外儲存維護一致性。 |
|
串流 |
事件持續到達、Agent 持續響應——Flink 流處理引擎為這種工作模式提供高吞吐、低延遲核心。 |
|
可信生產 |
基於分布式 Checkpoint,節點宕機自愈、狀態不丟、事件不丟、資料exactly-once。 |
Agent 類型
可基於不同情境構建 Workflow 或 ReAct Agent。
|
類型 |
適用情境 |
說明 |
|
Workflow Agent |
流程明確的情境 |
通過預定義事件驅動流程編排 Agent 行為。詳見Workflow Agent。 |
|
ReAct Agent |
需要靈活決策的情境 |
結合推理(Reasoning)與行動(Action),由 LLM 自主決定執行步驟。詳見ReAct Agent。 |
版本及模型說明
-
提交 Flink Agents 作業請選擇 VVR-11.7 及以上版本,否則作業因缺少運行時依賴而失敗。
VVR 引擎版本
Flink 版本
Flink Agents 版本
vvr-11.7.0-flink-1.20
1.20
0.2.1
-
模型由平台統一提供,無需單獨申請 API Key 即可在作業中調用。該功能目前處於白名單公測階段,詳情請參見Flink AI服務(內建模型)。
前提條件
-
產品與許可權
-
已開通Realtime Compute Flink 版並建立工作空間,詳見開通Realtime ComputeFlink版。
-
使用 RAM 使用者或 RAM 角色訪問時,已具備 Flink 控制台相關許可權,詳見許可權管理。
-
-
本地開發環境
-
Python:Python 3.10 或 3.11。
-
Java:Java 11+、Maven 3+。
-
快速入門
本節通過商品評價分析樣本示範完整流程:使用 Flink Agents 架構構建包含評價分析 Agent 的樣本Flink 作業,作業運行時,評價分析 Agent 調用 Flink 內建大模型服務,對評價文本輸出滿意度評分(1-5 分)和不滿意原因。
步驟一:下載樣本作業檔案
Python 作業
下載 quickstart-python.zip,包含以下檔案:
|
檔案 |
說明 |
|
main.py |
作業入口。建立執行環境、註冊 LLM 串連、構建流處理 pipeline。 |
|
review_analysis_agent.py |
Agent 實現。定義 Prompt 模板、ChatModel 配置和 Action 處理邏輯。 |
下載後的 ZIP 包直接用於上傳,無需解壓。樣本使用的 openai、dashscope 等 Python 依賴已預裝於 VVR 引擎鏡像。
Java 作業
下載以下檔案:
|
檔案 |
說明 |
|
測試 fat JAR,可直接上傳。 |
|
|
Java 原始碼,供參考。 |
步驟二:上傳並部署作業
-
單擊目標工作空間操作列的控制台。
-
在左側導覽列單擊檔案管理,單擊上傳資源,上傳下載的作業檔案。
-
在營運中心 > 作業營運頁面,單擊部署作業,按作業類型選擇 Python 作業或 JAR 作業,填寫部署資訊。
Python 作業部署配置
|
參數 |
樣本 |
|
部署模式 |
流模式 |
|
部署名稱 |
flink-agents-quickstart-python |
|
引擎版本 |
vvr-11.7.0-flink-1.20 |
|
Python 檔案地址 |
quickstart-python.zip |
|
Entry Module |
quickstart.main |
|
部署目標 |
default-queue |
Java 作業部署配置
需要添加額外的附加依賴檔案flink-agents-dist.jar。
|
參數 |
樣本 |
|
部署模式 |
流模式 |
|
部署名稱 |
flink-agents-quickstart-java |
|
引擎版本 |
vvr-11.7.0-flink-1.20 |
|
JAR URI |
quickstart-java.jar |
|
Entry Point Class |
org.apache.flink.agents.quickstart.Main |
|
附加依賴檔案 |
|
|
部署目標 |
default-queue |
運行參數
在部署詳情 > 運行參數配置 > 編輯 > 其他配置中添加 Flink 配置項。
-
Python 作業
python.executable: python3.10 python.client.executable: python3.10 containerized.master.env.FLINK_HOME: /flink containerized.taskmanager.env.FLINK_HOME: /flink classloader.parent-first-patterns.default: java.;scala.;com.esotericsoftware.kryo;org.apache.hadoop.;javax.annotation.;org.xml;javax.xml;org.apache.xerces;org.w3c;org.rocksdb.;org.slf4j;org.apache.log4j;org.apache.logging;org.apache.commons.logging;ch.qos.logback -
Java 作業
classloader.parent-first-patterns.default: java.;scala.;com.esotericsoftware.kryo;org.apache.hadoop.;javax.annotation.;org.xml;javax.xml;org.apache.xerces;org.w3c;org.rocksdb.;org.slf4j;org.apache.log4j;org.apache.logging;org.apache.commons.logging;ch.qos.logback配置項說明
配置項
說明
python.executable/python.client.executable指定 Python 解譯器版本(僅 Python 作業)。
containerized.{master,taskmanager}.env.FLINK_HOME設定 JobManager 與 TaskManager 的
FLINK_HOME環境變數(僅 Python 作業)。classloader.parent-first-patterns.default設定優先從父ClassLoader載入的包名首碼。
說明環境變數需同時配置
containerized.master.env.與containerized.taskmanager.env.兩個首碼,確保 JobManager 和 TaskManager 均能讀取。
步驟三:啟動作業並查看結果
-
在營運中心 > 作業營運頁面,單擊目標作業操作列的啟動。
-
在作業啟動對話方塊中選擇無狀態啟動,單擊啟動。
-
樣本使用記憶體資料來源,處理完成後作業自動結束。等待狀態變為已完成後,單擊作業名稱進入詳情。
-
在 TaskManager 頁簽查看日誌,搜尋
Review analysis result:關鍵字查看輸出結果。
生產環境通常使用 Kafka 等流式資料來源,作業持續運行。
開發自訂 Agent
安裝 Flink Agents
Python
建議使用虛擬環境:
python3 -m venv flink-agents-env
source flink-agents-env/bin/activate
pip install flink-agents
Java
在 pom.xml 中添加依賴:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-agents-api</artifactId>
<version>${flink-agents.version}</version>
<scope>provided</scope>
</dependency>
Flink Agents 相關依賴必須聲明 <scope>provided</scope>,由 VVR 引擎提供。
需要在本地 IDE 調試時,額外添加:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-agents-ide-support</artifactId>
<version>${flink-agents.version}</version>
<scope>provided</scope>
</dependency>
編寫 Agent
Flink Agents 作業的核心結構如下:
-
註冊 LLM 串連:通過
AgentsExecutionEnvironment註冊 ChatModel 串連,串連資訊(API Key、端點地址)通過ResourceDescriptor配置。詳見Chat Models。 -
定義 Agent:實現自訂 Agent,配置 Prompt 模板,通過
@action(Python)或@Action(Java)註解定義事件處理邏輯。 -
構建 Pipeline:將輸入資料流接入 Agent 處理,輸出分析結果。
Agent、Prompt、Tool、Memory 等概念詳見Flink Agents 開發文檔。
本地測試
本地測試由 Flink MiniCluster 自動啟動,無需部署叢集。
Python
python your_agent_job.py
Java
mvn exec:java -Dexec.mainClass="com.example.YourAgentJob"
依賴管理
鏡像預裝依賴
VVR 引擎鏡像預裝 Flink Agents 核心庫及部分依賴。
Python 側
|
依賴 |
說明 |
|
flink-agents 核心庫 |
Flink Agents Python API |
|
openai |
OpenAI 相容介面(含百鍊平台) |
|
dashscope |
阿里雲百鍊平台(通義千問原生介面) |
|
mcp |
Model Context Protocol |
鏡像未預裝的依賴需在提交作業時一併上傳。詳見使用Python依賴。
Java 側
VVR 引擎僅包含 Flink Agents thin JAR(核心代碼),所有 LLM integration 第三方依賴需要打入作業 JAR。
Java 作業依賴管理
使用 maven-shade-plugin 打包為 fat JAR,依賴按以下規則配置 scope:
-
Flink Agents 核心依賴(
flink-agents-api、flink-agents-runtime等):<scope>provided</scope>,不打入 fat JAR。 -
Flink 核心依賴(
flink-streaming-java等):<scope>provided</scope>。 -
LLM integration 依賴:按以下兩種情況處理。
-
社區 Flink Agents 已支援的 integration(如 OpenAI、Anthropic 等):
<scope>provided</scope>,不打入 fat JAR。社區在 Flink Agents 0.2.1 版本中支援的完整 integration 列表,請參見Built-in Providers。 -
社區未支援的 integration:
<scope>compile</scope>,打入 fat JAR。
-
-
與 Flink 衝突的庫(如 Jackson):在
maven-shade-plugin中配置 relocation。
完整 pom.xml 配置樣本請參見 quickstart-java-src.zip。
部署作業時,除業務 fat JAR 外,還需將 flink-agents-dist.jar 作為附加依賴一併上傳。該 JAR 包含 Flink Agents 運行時及社區已支援的全部 integration,自訂 Agent 作業同樣依賴該包。
提交作業到Realtime Compute Flink
Python 作業
提交方式與 PyFlink 作業一致,詳見Python作業開發。關鍵參數:
|
參數 |
說明 |
|
引擎版本 |
選擇支援 Flink Agents 的 VVR 引擎版本,詳見版本說明。 |
|
Python 檔案地址 |
作業入口 Python 檔案或 ZIP 包。 |
|
Entry Module |
ZIP 包入口模組名。 |
|
Python Libraries |
額外依賴包。 |
Java 作業
提交方式與 Flink JAR 作業一致,詳見JAR作業開發。關鍵參數:
|
參數 |
說明 |
|
引擎版本 |
選擇支援 Flink Agents 的 VVR 引擎版本,詳見版本說明。 |
|
JAR URI |
fat JAR 路徑。 |
|
Entry Point Class |
作業入口類全限定名。 |
通用運行參數
Python 解譯器版本
VVR 引擎鏡像預設 Python 版本為 3.9,Flink Agents 需要 3.10 或 3.11。Python 作業必須配置:
python.executable: python3.10
python.client.executable: python3.10
containerized.master.env.FLINK_HOME: /flink
containerized.taskmanager.env.FLINK_HOME: /flink
自訂環境變數
作業代碼通過環境變數讀取配置(如端點地址,切換模型)時,需同時配置 master 與 taskmanager 兩個首碼:
containerized.master.env.<ENV_VAR_NAME>: <value>
containerized.taskmanager.env.<ENV_VAR_NAME>: <value>
切換模型
containerized.master.env.OPENAI_MODEL: qwen3.6-flash
containerized.taskmanager.env.OPENAI_MODEL: qwen3.6-flash