本文介绍如何在 Flink SQL 作业中使用内置大模型服务进行流式情感分析与向量化,无需自行申请 API Key。
背景信息
Flink AI 服务提供托管式内置大模型能力,用户无需自行申请 API Key,即可在 Flink SQL 作业中直接调用内置模型,实现流式 AI 推理与向量化。以下是两种核心模型的应用场景:
-
chat/completions模型 chat/completions模型是一种基于对话生成和文本理解的大模型,广泛应用于情感分析、意图识别、问答系统等场景。
-
情感分析:对企业社交媒体评论进行实时情感分类,快速识别用户情绪(如正面、负面、中性)。
-
智能客服:通过对话生成能力,为用户提供自然语言交互的智能客服服务。
-
内容审核:自动检测文本中的敏感内容或违规信息,提升内容安全审核效率。
-
-
embedding模型 embedding模型能够将文本转换为高维向量表示,适用于语义搜索、推荐系统、知识图谱构建等场景。
-
语义搜索:通过对商品描述或用户查询进行向量化处理,实现基于语义的相关性搜索。
-
推荐系统:利用文本向量化技术,挖掘用户兴趣与商品特征之间的潜在关联,提升推荐精准度。
-
知识图谱:将非结构化文本转化为向量形式,便于后续的知识抽取和关系建模。
-
前提条件
-
已开通Flink工作空间,详情请参见开通实时计算Flink版。
-
已开通 Flink AI服务(内置模型)。
使用限制
-
仅实时计算引擎 VVR 11.7 及以上版本支持。
-
ML_PREDICT 算子的吞吐量受模型平台限流限制,触及流量上限时作业可能出现反压或超时重启。
步骤一:注册内置模型
请参见模型设置注册 Flink 内置模型。
chat/completions模型任务
以注册内置文本模型为例,SQL代码如下:
CREATE MODEL ai_analyze_sentiment
INPUT (`input` STRING)
OUTPUT (`content` STRING)
WITH (
'provider' = 'dashscope',
'task' = 'chat/completions',
'model' = 'qwen3.5-flash',
'system-prompt' = 'Classify the text below into one of the following labels: [positive, negative, neutral, mixed]. Output only the label.'
);
embedding模型任务
以注册内置向量模型为例,SQL代码如下:
CREATE MODEL embedding_model
INPUT (`input` STRING)
OUTPUT (`embeddings` ARRAY<FLOAT>)
WITH (
'provider' = 'dashscope',
'task' = 'embeddings',
'model' = 'text-embedding-v4'
);
步骤二:创建作业
请参见Flink SQL作业创建SQL流作业草稿。
步骤三:编写SQL作业进行AI大模型分析
chat/completions模型任务
使用AI函数通用调用已注册的ai_analyze_sentiment模型,对电影评论进行情感分析。
拷贝如下示例SQL到SQL编辑区域。
--创建临时结果表
CREATE TEMPORARY TABLE print_sink(
id BIGINT,
movie_name VARCHAR,
predict_label VARCHAR,
actual_label VARCHAR
) WITH (
'connector' = 'print', -- print连接器
'logger' = 'true' -- 控制台显示计算结果
);
-- 创建临时数据视图构造测试数据
-- | id | movie_name | comment | actual_label |
-- | 1 | 好东西 | 最爱小孩子猜声音那段,算得上看过的电影里相当浪漫的叙事了。很温和也很有爱。| POSITIVE |
-- | 2 | 水饺皇后 | 乏善可陈 | NEGATIVE |
CREATE TEMPORARY VIEW movie_comment(id, movie_name, user_comment, actual_label)
AS VALUES (1, '好东西', '最爱小孩子猜声音那段,算得上看过的电影里相当浪漫的叙事了。很温和也很有爱。', 'positive'), (2, '水饺皇后', '乏善可陈', 'negative');
INSERT INTO print_sink
SELECT id, movie_name, content as predict_label, actual_label
FROM ML_PREDICT(
TABLE movie_comment,
MODEL ai_analyze_sentiment, -- 注册的 Flink 内置文本模型
DESCRIPTOR(user_comment));
embedding模型任务
使用AI函数通用调用已注册的embedding_model模型,对电影评论进行情感分析,并将结果写入Milvus(公测中)。
拷贝如下示例SQL到SQL编辑区域。
-- 创建临时结果表milvus_sink
CREATE TEMPORARY TABLE milvus_sink
(
id STRING,
movie_name STRING,
user_comment STRING,
embeddings ARRAY<FLOAT>,
PRIMARY KEY (id) NOT ENFORCED
)
WITH (
'connector' = 'milvus',
'endpoint' = '<YOUR-ENDPOINT>',
'port' = '<YOUR-PORT>',
'userName' = '<YOUR-USERNAME>',
'password' = '<YOUR-PASSWORD>',
'databaseName' = 'default',
'collectionName' = 'movie-comment-embeddings'
);
-- 创建临时数据视图构造测试数据
-- | id | movie_name | comment | actual_label|
-- | 1 | 好东西 |最爱小孩子猜声音那段,算得上看过的电影里相当浪漫的叙事了。很温和也很有爱。| POSITIVE |
-- | 2 | 水饺皇后 | 乏善可陈 | NEGATIVE |
CREATE TEMPORARY VIEW movie_comment(id, movie_name, user_comment)
AS VALUES ('1', '好东西', '最爱小孩子猜声音那段,算得上看过的电影里相当浪漫的叙事了。很温和也很有爱。'), ('2', '水饺皇后', '乏善可陈');
INSERT INTO
milvus_sink
SELECT
id,
movie_name,
user_comment,
embeddings
FROM
ML_PREDICT (
TABLE movie_comment,
MODEL embedding_model, -- 注册的 Flink 内置向量模型
DESCRIPTOR (user_comment)
);
步骤四:作业部署并启动
请参见Flink SQL作业部署作业并启动。
步骤五:查看分析结果
chat/completions模型任务
-
查看目标作业状态已完成。
-
在页面,单击目标作业名称。
-
在作业日志页签,选择Task Managers页签下的当前TaskManager。
-
单击日志,在页面搜索PrintSinkOutputWriter相关的日志信息。
经模型分析后的结果
predict_label和实际结果actual_label一致。选择 Task Managers 页签,在左侧选择运行日志,日志中可看到 PrintSinkOutputWriter 输出结果,例如
+I[1, 好东西, positive, positive]和+I[2, 水饺皇后, negative, negative],表明 predict_label 与 actual_label 一致。
相关文档
-
AI模型DDL数据定义语句请参见模型设置。
-
AI函数请参见通用调用。
-
向量检索服务Milvus(公测中)。