Flink SQL调用Agent(公测)

更新时间:
复制 MD 格式

实时计算Flink版提供AGENT_RUN表值函数(TVF),支持在Flink SQL作业中调用基于Flink Agents框架构建的Agent。本文介绍AGENT_RUN的语法、参数、使用示例及作业提交方式。

前提条件

语法

SELECT * FROM 
AGENT_RUN(
   TABLE input_table [PARTITION BY key]
  ,AGENT agent_name
 [,LANGUAGE 'java'/'python']
 [,CONFIG => MAP['key', 'value']]
)

参数说明

参数

是否必填

说明

input_table

包含待处理数据的输入表。

key

输入表的分区键。指定后,相同key的数据由同一Agent实例处理。

agent_name

Agent的全限定类名。需为通过Flink Agents框架构建的Agent。

LANGUAGE

Agent的开发语言,取值为javapython。默认为java

CONFIG

Flink Agents配置参数,以MAP形式传入键值对。详情请参见配置参数

输出说明

输出表包含输入表的所有列,以及Agent的输出列。输出列固定为output VARIANT类型。

使用示例

Java Agent

CREATE TEMPORARY TABLE reviews
(
    id   INT
    ,msg STRING
)
WITH (
    'connector' = 'datagen'
    ,'number-of-rows' = '3'
    ,'fields.id.kind' = 'sequence'
    ,'fields.id.start' = '1'
    ,'fields.id.end' = '3'
    ,'fields.msg.length' = '4'
);

CREATE TEMPORARY TABLE `test`
(
    id      INT,
    output VARIANT
)
WITH ('connector' = 'print');

INSERT INTO `test`
SELECT
    id,
    output
FROM AGENT_RUN(TABLE reviews PARTITION BY id,'org.apache.flink.agents.examples.agents.EchoAgent');

Python Agent

CREATE TEMPORARY TABLE reviews
(
    id   INT
    ,msg STRING
)
WITH (
    'connector' = 'datagen'
    ,'number-of-rows' = '3'
    ,'fields.id.kind' = 'sequence'
    ,'fields.id.start' = '1'
    ,'fields.id.end' = '3'
    ,'fields.msg.length' = '4'
);

CREATE TEMPORARY TABLE `test`
(
    id      INT
    ,output VARIANT
)
WITH ('connector' = 'print');

INSERT INTO `test`
SELECT
    id
    ,output
FROM AGENT_RUN(TABLE reviews PARTITION BY id
    ,'echo_agent.SqlEchoPythonAgent'
    ,'python');

作业依赖配置

Java Agent作业

将包含Agent类的JAR包通过附加依赖上传,然后提交SQL作业即可。

Python Agent作业

  1. 将包含Python Agent类的Python文件通过附加依赖上传。

  2. 将包含Flink Agents Python Wheel和作业相关依赖的虚拟环境(venv)上传。构建方式请参见使用自定义的Python虚拟环境

  3. 部署详情 > 运行参数配置 > 其他配置中添加以下配置:

参数

说明

python.files

通过附加依赖上传的Python文件路径,格式为oss://${bucket}/<文件名>.py。在运维界面可查看其OSS路径。

python.archives

包含Flink Agents及作业相关依赖的venv压缩包路径。

python.client.executable

上传的venvPython解释器路径,例如venv.zip/venv/bin/python

python.executable

固定为python3.10

containerized.master.env.FLINK_HOME

固定为/flink

containerized.taskmanager.env.FLINK_HOME

固定为/flink

配置示例如下:

python.files: 'oss://${bucket}/echo_agent.py'
python.archives: >-oss://${bucket}/venv.zip
python.executable: python3.10
python.client.executable: venv.zip/venv/bin/python
containerized.master.env.FLINK_HOME: /flink
containerized.taskmanager.env.FLINK_HOME: /flink