基于YAML CDC实现AI智能文本标注

更新时间:
复制 MD 格式

本文介绍如何通过Flink CDC YAML作业结合AI Function,在数据入湖过程中对非结构化文本进行实时智能标注。

方案架构

企业每天产生大量非结构化文本数据(客服工单、用户评论、投诉记录等),传统方式依赖人工逐条标注或T+1离线批处理,成本高、延迟大、标准不一致。

通过CDC YAMLTransform模块调用AI Function,数据在从Kafka流入Paimon数据湖的过程中同步完成智能标注,实现秒级延迟的流式打标。

image

该方案涵盖以下典型场景:

  • 客诉风险智能标注:使用AI_COMPLETE对客服工单进行投诉情绪等级判定、曝光意图识别等多维度风险打标。

  • 内容分类智能标注:使用AI_CLASSIFY对用户评论或反馈进行预定义类别的自动分类。

  • 敏感信息自动脱敏:使用AI_MASK对入湖数据中的个人隐私信息(PII)进行识别和脱敏。

前提条件

使用限制

  • AI FunctionCDC YAML中仅支持在Transform模块的projection表达式中使用。

  • AI Function的吞吐量受阿里云百炼平台限流限制。触及流量上限时,Flink作业可能出现反压或超时重启。详情请参见百炼平台限流

当前版本支持的内置百炼模型如下:

模型名称

说明

qwen3.6-flash

通用文本生成,推荐用于标注场景

qwen3.6-plus

更强推理能力

qwen3.5-flash

通用文本生成

qwen3.5-plus

更强推理能力

text-embedding-v4

文本向量化

AI Function速查

CDC YAMLTransform模块支持以下AI Function。

函数

签名

返回字段

用途

AI_COMPLETE

AI_COMPLETE('modelName', input, systemPrompt)

result

通用文本生成

AI_CLASSIFY

AI_CLASSIFY('modelName', input, labels)

categoryconfidence

文本分类

AI_TRANSLATE

AI_TRANSLATE('modelName', input, sourceLang, targetLang)

translated_textdetected_language

文本翻译

AI_SUMMARIZE

AI_SUMMARIZE('modelName', input, maxLength)

summary

文本摘要

AI_SENTIMENT

AI_SENTIMENT('modelName', input)

scorelabelconfidence

情感分析

AI_EXTRACT

AI_EXTRACT('modelName', input, schema)

extracted_json

结构化信息提取

AI_MASK

AI_MASK('modelName', input, entities)

masked_textdetected_entities

敏感信息脱敏

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格式访问结构化内容。如果下游需要特定类型(如STRINGINT),需通过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=()
}

部署与启动

  1. 登录实时计算控制台,创建CDC YAML数据摄入作业,将上述任一场景的YAML配置粘贴到编辑器中。详情请参见Flink CDC数据摄入作业开发

  2. 上传附加依赖文件。在作业草稿的附加依赖文件区域,上传以下JAR包:

    flink-cdc-pipeline-model-openai-compatible

  3. 部署并启动作业。

监控与限流

当前版本仅通过Flink Web UI监控AI FunctionToken用量。在 Flink Dashboard 中,单击目标作业,选择 Metrics 页签。在 Add Metric 下拉列表中选择目标算子的指标(例如 ai_model_completion_tokens),查看当前指标值。