本文介绍如何在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
在您使用 Flink AI 服务时,连接地址和访问凭证由平台内置并动态提供,无需在作业中填写 endpoint 和 api-key。
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 服务用于推理的实际模型名称。提供的模型情况请参见内置模型列表。 |
endpoint | STRING | 否 | 无 | OpenAI 兼容 API 服务的 Base URL;使用 Flink AI 服务时无需填写。 |
api-key | STRING | 否 | 无 | OpenAI 兼容 API 服务鉴权所用的 API Key;使用 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 AI 服务内置模型时无需配置 endpoint 和 api-key;连接外部 OpenAI 相容的 API 服务时,需要同时配置这两个参数。
复用 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