全部產品
Search
文件中心

Realtime Compute for Apache Flink:Flink Agents開發(公測)

更新時間:Jul 29, 2026

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服務(內建模型)。

前提條件

  • 產品與許可權

  • 本地開發環境

    • 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 作業

下載以下檔案:

檔案

說明

quickstart-java.jar

測試 fat JAR,可直接上傳。

quickstart-java-src.zip

Java 原始碼,供參考。

步驟二:上傳並部署作業

  1. 登入Realtime Compute控制台。

  2. 單擊目標工作空間操作列的控制台。

  3. 在左側導覽列單擊檔案管理,單擊上傳資源,上傳下載的作業檔案。

  4. 在營運中心 > 作業營運頁面,單擊部署作業,按作業類型選擇 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

附加依賴檔案

flink-agents-dist.jar

部署目標

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 均能讀取。

步驟三:啟動作業並查看結果

  1. 在營運中心 > 作業營運頁面,單擊目標作業操作列的啟動。

  2. 在作業啟動對話方塊中選擇無狀態啟動,單擊啟動。

  3. 樣本使用記憶體資料來源,處理完成後作業自動結束。等待狀態變為已完成後,單擊作業名稱進入詳情。

  4. 在 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       

詳見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 作業的核心結構如下:

  1. 註冊 LLM 串連:通過 AgentsExecutionEnvironment 註冊 ChatModel 串連,串連資訊(API Key、端點地址)通過 ResourceDescriptor 配置。詳見Chat Models。

  2. 定義 Agent:實現自訂 Agent,配置 Prompt 模板,通過 @action(Python)或 @Action(Java)註解定義事件處理邏輯。

  3. 構建 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"       

詳見Flink Agents 部署文檔。

依賴管理

鏡像預裝依賴

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