实时计算Flink版提供AGENT_RUN表值函数(TVF),支持在Flink SQL作业中调用基于Flink Agents框架构建的Agent。本文介绍AGENT_RUN的语法、参数、使用示例及作业提交方式。
前提条件
已基于Flink Agents框架构建Agent,详情请参见Flink Agents开发(公测)。
使用实时计算引擎VVR 11.8及以上版本。
语法
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的开发语言,取值为 |
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作业
将包含Python Agent类的Python文件通过附加依赖上传。
将包含Flink Agents Python Wheel和作业相关依赖的虚拟环境(venv)上传。构建方式请参见使用自定义的Python虚拟环境。
在中添加以下配置:
参数 | 说明 |
python.files | 通过附加依赖上传的Python文件路径,格式为 |
python.archives | 包含Flink Agents及作业相关依赖的venv压缩包路径。 |
python.client.executable | 上传的venv中Python解释器路径,例如 |
python.executable | 固定为 |
containerized.master.env.FLINK_HOME | 固定为 |
containerized.taskmanager.env.FLINK_HOME | 固定为 |
配置示例如下:
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