本文介绍如何在Flink数据摄入作业中使用内置的 AI 推理模型,实现对数据流的实时智能处理(如文本摘要生成、向量化等)。
前提条件
-
实时计算引擎:VVR 11.9(含 Preview 预览版本)及更高版本。
-
附加依赖包:将 flink-cdc-pipeline-model-openai-compatible JAR 包作为作业附加依赖上传。
工作原理
AI 推理模型的集成分为两个步骤:
-
定义模型:在 YAML 作业草稿的 Pipeline 块中定义模型,也可以直接引用元数据中心中创建的模型。
-
使用模型:在 YAML 作业草稿的 Transform 块中调用 AI 推理或向量生成模型。
每个模型需要指定唯一名称,该名称将作为函数名在数据转换中调用。
配置并调用 Flink AI 内置模型
步骤一:先决条件准备
开始前,请先参照 Flink AI服务(内置模型)文档的说明开通 Flink AI 服务。
此外,您还需要将 flink-cdc-pipeline-model-openai-compatible JAR 包作为作业附加依赖上传,并附加到作业草稿中。
步骤二:在 Pipeline 块中定义模型
为了定义 AI 模型,您需要在 YAML 作业草稿的 pipeline 模块中添加 model 字段列表。您需要为每个模型提供唯一的名称,且需要符合一般标识符的规则(必须以字母开头,且只能包含字母、下划线及数字)。
必须填写的参数包括:
-
模型标识符
name,可任意选择,需要为合法的标识符 -
模型类型
type,固定为openai-compatible -
调用百炼服务中的模型名称
model
pipeline:
model:
- name: CHAT
type: openai-compatible
model: qwen-plus
user-prompt: Please summarize the input.
- name: GET_EMBEDDING
type: openai-compatible
model: text-embedding-v4
以下是您可以指定的完整参数列表:
|
参数 |
类型 |
是否必填 |
默认值 |
说明 |
|
name |
STRING |
是 |
无 |
Pipeline 中的模型别名,供 AI 函数引用;必须唯一,且符合标识符命名规则。 |
|
type |
STRING |
是 |
无 |
模型类型,这里固定填写 openai-compatible。 |
|
model |
STRING |
是 |
无 |
Flink AI 服务用于推理的实际模型名称。提供的模型情况请参见内置模型列表。 |
|
task |
STRING |
否 |
cdc-yaml-ai-task |
模型任务名称,作为监控指标中的 task 分组。 |
|
error-handling-strategy |
STRING |
否 |
retry |
请求失败后的处理策略;可选值:retry、failover、ignore。 |
|
retry-num |
INT |
否 |
100 |
当 error-handling-strategy 为 retry 时,单次模型调用的最大尝试次数。 |
|
retry-fallback-strategy |
STRING |
否 |
failover |
重试耗尽后的处理策略;可选值:failover、ignore,不能配置为 retry。 |
|
retry-backoff-strategy |
STRING |
否 |
fixed |
重试间隔策略;可选值:fixed、exponential。 |
|
retry-backoff-base-interval |
DURATION |
否 |
1 s |
重试基础间隔,例如 500 ms、1 s、2 min。 |
|
system-prompt |
STRING |
否 |
无 |
添加到请求中的系统提示词;常规文本生成时会放在 AI 函数自身的系统提示词之前。 |
|
user-prompt |
STRING |
否 |
无 |
在模型输入之后追加一条用户提示词。 |
|
temperature |
DOUBLE |
否 |
由模型服务决定 |
采样温度。 |
|
top-p |
DOUBLE |
否 |
由模型服务决定 |
核采样概率质量。 |
|
stop |
STRING |
否 |
无 |
遇到指定字符串时停止生成。 |
|
max-tokens |
INT |
否 |
由模型服务决定 |
单次请求允许生成的最大 Token 数。 |
|
presence-penalty |
DOUBLE |
否 |
由模型服务决定 |
Presence penalty,用于降低模型重复已有主题的倾向。 |
|
n |
LONG |
否 |
由模型服务决定 |
请求生成的候选结果数量;当前客户端只返回第一条结果。 |
|
seed |
LONG |
否 |
无 |
随机种子;模型服务支持时可用于提高结果的可复现性。 |
|
response-format |
STRING |
否 |
json_object |
响应格式;当前只支持 json_object。 |
|
content-type |
STRING |
否 |
text |
模型输入类型;可选值:text、image_url。image_url 会将输入字符串作为图片 URL 发送。 |
|
extra-header |
STRING |
否 |
无 |
以 JSON 对象形式声明附加到请求中的 HTTP Header。 |
|
extra-body |
STRING |
否 |
无 |
以 JSON 对象形式声明附加到请求体中的厂商扩展参数。 |
|
dimension |
INT |
否 |
由模型服务决定 |
AI_EMBED 返回向量的维度;仅在对应模型支持自定义维度时生效。 |
步骤三:在 Transform 模块中调用 AI Function
按照模态分类,下面是您可以在 Transform 模块中调用的内置 AI 函数。
每个内置 AI 函数的第一个参数均为您在 Pipeline 块中定义的模型标识符(model 中 name 参数的值)。
文本生成类函数
-
AI_COMPLETE
-
用途:通用的文本生成函数,可根据给出的提示词处理输入文本,适用于内容生成、文本改写和通用问答。
-
使用示例:
AI_COMPLETE('chat_model', content_string, '提示:请概括输入内容') -
返回示例:
{"result":"输入内容的摘要"}
-
-
AI_CLASSIFY
-
用途:将输入文本分类到指定的候选类别中。多个类别之间使用逗号分隔。
-
使用示例:
AI_CLASSIFY('chat_model', content_string, '新闻,体育,娱乐') -
返回示例:
{"category":"体育","confidence":0.95}
-
-
AI_TRANSLATE
-
用途:将输入文本从源语言翻译为目标语言。源语言填写
auto时,由模型自动识别。 -
使用示例:
AI_TRANSLATE('chat_model', content_string, 'auto', 'en') -
返回示例:
{"translated_text":"Hello, world.","detected_language":"zh"}
-
-
AI_SUMMARIZE
-
用途:为输入文本生成摘要。第三个参数表示摘要允许包含的最大字符数。
-
使用示例:
AI_SUMMARIZE('chat_model', content_string, 200) -
返回示例:
{"summary":"输入内容的摘要"}
-
-
AI_SENTIMENT
-
用途:分析输入文本的情感倾向。情感分数
score的取值范围为-1.0~1.0。 -
使用示例:
AI_SENTIMENT('chat_model', review_content) -
返回示例:
{"score":0.8,"label":"positive","confidence":0.92}
-
-
AI_EXTRACT
-
用途:按照指定的 Schema 从非结构化文本中提取结构化信息。
-
使用示例:
AI_EXTRACT('chat_model', content_string, '{"name":"string","age":"integer"}') -
返回示例:
{"extracted_json":"{\"name\":\"张三\",\"age\":25}"}
-
-
AI_MASK
-
用途:识别并遮盖指定类型的敏感信息。多个实体类型之间使用逗号分隔。
-
使用示例:
AI_MASK('chat_model', content_string, 'name,phone,email') -
返回示例:
{"masked_text":"联系人:张*,手机:138****8000","detected_entities":"name,phone"}
-
文本向量化函数
-
AI_EMBED
-
用途:将输入文本转换为浮点数向量,可用于语义检索、相似度计算和文本聚类。
-
使用示例:
AI_EMBED('embedding_model', content_string) -
返回示例:
[0.0123, -0.0876, 0.2345, ...]
-
多模态生成类函数
-
AI_IMAGE_COMPLETE
-
用途:根据图片二进制数据和提示词生成文本,可用于图片描述、内容识别和视觉问答。
-
使用示例:
AI_IMAGE_COMPLETE('vision_model', image_bytes, '请描述图片中的商品') -
返回示例:
"图片中是一双白色运动鞋。"
-
多模态向量化函数
-
AI_IMAGE_EMBED
-
用途:将图片二进制数据转换为浮点数向量,可用于以图搜图、图片聚类和多模态检索。
-
使用示例:
AI_IMAGE_EMBED('image_embedding_model', image_bytes) -
返回示例:
[0.0312, -0.1024, 0.2867, ...]
-
完整配置示例
以下示例从 MySQL 读取文章及图片数据,调用文本摘要和向量化函数,并写入 DLF Paimon 数据湖。
默认情况下,AI 推理模型会顺序串行求值,受限于网络 I/O 速率可能导致吞吐下降、Checkpoint 过慢等问题。
您可以开启 transform.async-execution.enabled参数启用异步 AI 推理优化(实验性功能),具体的参数功能及限制参见 Pipeline配置参数。
source:
type: mysql
using.built-in-catalog: mysql_catalog
tables: inventory.articles
sink:
type: paimon
using.built-in-catalog: dlf_catalog
transform:
- source-table: inventory.articles
projection: >-
id,
AI_SUMMARIZE('chat_model', content, 200) AS summary,
AI_EMBED('embedding_model', content) AS text_embedding,
AI_IMAGE_COMPLETE('chat_model', image_bytes, '请描述图片内容') AS image_description,
AI_IMAGE_EMBED('embedding_model', image_bytes) AS image_embedding
pipeline:
# 启用异步 AI 推理优化试验性特性
transform.async-execution.enabled: true
parallelism: 4
model:
- name: chat_model
type: openai-compatible
model: qwen3.7-max
- name: embedding_model
type: openai-compatible
model: qwen3-vl-embedding
dimension: 1024
复用 Flink SQL 的 AI 模型定义
目前只支持 Flink SQL 定义的 openai-compat、dashscope 类型的 AI 模型。
创建推理模型
VVR 11.7 及以上版本支持通过 CREATE MODEL 在元数据中心创建持久化模型。使用内置模型时,无需配置 endpoint 和 api-key。
CREATE MODEL qwen_chat
USING openai_compatible
WITH (
'provider' = 'openai-compat',
'task' = 'chat/completions',
'model' = 'qwen3.7-max'
);
CREATE MODEL multimodal_embedding
USING openai_compatible
WITH (
'provider' = 'dashscope',
'task' = 'multimodal-embedding',
'model' = 'qwen3-vl-embedding',
'content-type' = 'image_url',
'dimension' = '512'
);
复用已有模型
创建完成后,可以在 YAML Pipeline 中通过完整名称引用 Flink SQL 模型:
pipeline:
model:
- name: chat_model
reuse-model: vvp.default.qwen_chat
- name: image_embedding_model
reuse-model: vvp.default.multimodal_embedding