本文介绍如何通过Flink CDC YAML作业结合AI Function,在数据入湖过程中对非结构化文本进行实时智能标注。
方案架构
企业每天产生大量非结构化文本数据(客服工单、用户评论、投诉记录等),传统方式依赖人工逐条标注或T+1离线批处理,成本高、延迟大、标准不一致。
通过CDC YAML的Transform模块调用AI Function,数据在从Kafka流入Paimon数据湖的过程中同步完成智能标注,实现秒级延迟的流式打标。
该方案涵盖以下典型场景:
-
客诉风险智能标注:使用
AI_COMPLETE对客服工单进行投诉情绪等级判定、曝光意图识别等多维度风险打标。 -
内容分类智能标注:使用
AI_CLASSIFY对用户评论或反馈进行预定义类别的自动分类。 -
敏感信息自动脱敏:使用
AI_MASK对入湖数据中的个人隐私信息(PII)进行识别和脱敏。
前提条件
-
已开通实时计算Flink版并创建工作空间,详情请参见开通实时计算Flink版。
-
已开通Flink AI服务(内置模型)。
-
仅实时计算引擎版本为VVR 11.9(含preview)及以上支持。
-
已创建Kafka实例,并准备好待标注的文本数据Topic。
使用限制
-
AI Function在CDC YAML中仅支持在Transform模块的
projection表达式中使用。 -
AI Function的吞吐量受阿里云百炼平台限流限制。触及流量上限时,Flink作业可能出现反压或超时重启。详情请参见百炼平台限流。
当前版本支持的内置百炼模型如下:
|
模型名称 |
说明 |
|
qwen3.6-flash |
通用文本生成,推荐用于标注场景 |
|
qwen3.6-plus |
更强推理能力 |
|
qwen3.5-flash |
通用文本生成 |
|
qwen3.5-plus |
更强推理能力 |
|
text-embedding-v4 |
文本向量化 |
AI Function速查
CDC YAML的Transform模块支持以下AI Function。
|
函数 |
签名 |
返回字段 |
用途 |
|
AI_COMPLETE |
|
|
通用文本生成 |
|
|
|
文本分类 |
|
|
|
|
文本翻译 |
|
|
|
|
文本摘要 |
|
|
|
|
情感分析 |
|
|
|
|
结构化信息提取 |
|
|
|
|
敏感信息脱敏 |
在projection表达式中,通过['字段名']访问返回字段。例如:AI_COMPLETE('mymodel', content, '')['result']。
场景一:客诉风险智能标注
使用AI_COMPLETE对客服工单文本进行多维度风险打标,输出投诉情绪等级、曝光意图、品牌抹黑倾向等判定结果。
作业示例
部署时需要额外的依赖文件与配置,详情请参见部署与启动。
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: <YOUR-KAFKA-BROKER>:9092
topic: customer_complaints
value.format: json
scan.startup.mode: earliest-offset
# 将测试结果输出到控制台日志中
sink:
type: values
# sink:
# type: paimon
# catalog.properties.metastore: rest
# catalog.properties.uri: <YOUR-DLF-ENDPOINT>
# catalog.properties.warehouse: <YOUR-CATALOG-NAME>
# catalog.properties.token.provider: dlf
# table.properties.deletion-vectors.enabled: true
transform:
- source-table: customer_complaints
projection: >
comment,
AI_COMPLETE('mymodel', comment, '')['result'] as risk_label
pipeline:
parallelism: 2
model:
- name: mymodel
type: openai-compatible
model-name: qwen3.6-flash
system-prompt: '提示词:
#你是电商客诉风险判定专家,根据用户评价内容进行客观判标。
#判标规则
1.投诉情绪等级(三选一):正常/不满/愤怒
2.是否有明确曝光意图(二选一):是/否
3.是否有抹黑品牌倾向(二选一):是/否
#输出(严格按以下结构,一行一个)
投诉情绪等级:xxx
是否有明确曝光意图:xxx
是否有抹黑品牌倾向:xxx
置信度:0.xx
打标依据:xxx
'
输入输出示例
输入数据(JSON格式):
{
"comment": "买了不到一个月就坏了,联系客服一直让等,等了两周没人管。再不处理我,就发到网上曝光投诉你们。"
}
输出数据:
DataChangeEvent{
tableId=customer_complaints,
after=[
用户来电反映,2024年3月购买的某品牌空调...,
"投诉情绪等级:愤怒
是否有明确曝光意图:是
是否有抹黑品牌倾向:是
置信度:0.95
打标依据:再不处理我,就发到网上曝光投诉你们。"
],
op=INSERT, meta=()
}
AI Function输出字段的数据类型为VARIANT,便于以JSON格式访问结构化内容。如果下游需要特定类型(如STRING或INT),需通过CAST显式转换,例如CAST(risk_label AS STRING)。
场景二:内容分类智能标注
使用AI_CLASSIFY对用户评论或反馈文本进行预定义类别分类。
作业示例
部署时需要额外的依赖文件与配置,详情请参见部署与启动。
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: <YOUR-KAFKA-BROKER>:9092
topic: user_reviews
value.format: json
scan.startup.mode: earliest-offset
# 将测试结果输出到控制台日志中
sink:
type: values
# sink:
# type: paimon
# catalog.properties.metastore: rest
# catalog.properties.uri: <YOUR-DLF-ENDPOINT>
# catalog.properties.warehouse: <YOUR-CATALOG-NAME>
# catalog.properties.token.provider: dlf
# table.properties.deletion-vectors.enabled: true
transform:
- source-table: user_reviews
projection: >
content,
AI_CLASSIFY('mymodel', content, 'ARRAY[''产品质量'',''售后服务'',''物流问题'',''价格异议'']')['category'] as category,
AI_CLASSIFY('mymodel', content, 'ARRAY[''产品质量'',''售后服务'',''物流问题'',''价格异议'']')['confidence'] as confidence
pipeline:
parallelism: 2
model:
- name: mymodel
type: openai-compatible
model-name: qwen3.6-flash
输入输出示例
输入数据:
{"content": "空调噪音太大了,开机就嗡嗡响,完全没法睡觉"}
输出数据:
DataChangeEvent{
tableId=user_reviews,
before=[],
after=[空调噪音太大了,开机就嗡嗡响,完全没法睡觉, "产品质量", 0.95],
op=INSERT, meta=()
}
场景三:敏感信息自动脱敏
使用AI_MASK对入湖数据中的个人隐私信息(姓名、手机号、地址等)进行自动识别和脱敏。
作业示例
部署时需要额外的依赖文件与配置,详情请参见部署与启动。
source:
type: kafka
name: Kafka Source
properties.bootstrap.servers: <YOUR-KAFKA-BROKER>:9092
topic: user_feedbacks
value.format: json
scan.startup.mode: earliest-offset
# 将测试结果输出到控制台日志中
sink:
type: values
# sink:
# type: paimon
# catalog.properties.metastore: rest
# catalog.properties.uri: <YOUR-DLF-ENDPOINT>
# catalog.properties.warehouse: <YOUR-CATALOG-NAME>
# catalog.properties.token.provider: dlf
# table.properties.deletion-vectors.enabled: true
transform:
- source-table: user_feedbacks
projection: >
AI_MASK('mymodel', content, 'ARRAY[''姓名'',''手机号'',''身份证号'',''地址'']')['masked_text'] as masked_content,
AI_MASK('mymodel', content, 'ARRAY[''姓名'',''手机号'',''身份证号'',''地址'']')['detected_entities'] as detected_entities
pipeline:
parallelism: 2
model:
- name: mymodel
type: openai-compatible
model-name: qwen3.6-flash
输入输出示例
输入数据:
{"content": "客户张三,手机号13812345678,地址为北京市朝阳区建国路88号,反映冰箱不制冷"}
输出数据:
DataChangeEvent{
tableId=user_feedbacks,
before=[],
after=["客户张*,手机号138****5678,地址为北京市朝阳区**,反映冰箱不制冷", [
{"type":"姓名","value":"张三"},
{"type":"手机号","value":"13812345678"},
{"type":"地址","value":"北京市朝阳区建国路88号"}
]],
op=INSERT, meta=()
}
部署与启动
-
登录实时计算控制台,创建CDC YAML数据摄入作业,将上述任一场景的YAML配置粘贴到编辑器中。详情请参见Flink CDC数据摄入作业开发。
-
上传附加依赖文件。在作业草稿的附加依赖文件区域,上传以下JAR包:
-
部署并启动作业。
监控与限流
当前版本仅通过Flink Web UI监控AI Function的Token用量。在 Flink Dashboard 中,单击目标作业,选择 Metrics 页签。在 Add Metric 下拉列表中选择目标算子的指标(例如 ai_model_completion_tokens),查看当前指标值。