You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka SQL Server Sink Connector读取Topic数据失败求助

问题分析与解决

核心问题

错误日志已明确说明:当启用delete.enabled=true且pk.mode=record_key时,Kafka Sink Connector要求消息的Key必须包含非空的Schema信息(Struct或原始类型Schema),但当前Topic中的Key是纯字符串且无Schema,同时Sink配置里key.converter.schemas.enable=false,导致Connector无法识别Key的结构,最终抛出异常终止任务。

解决方案

方案1:关闭删除功能(最简单,若无需同步删除操作)

修改Sink Connector配置,将delete.enabled设为false即可绕过Key的Schema校验:

CREATE SINK CONNECTOR hsr_user_sqlsrvr1
WITH (
  'connector.class' = 'io.confluent.connect.jdbc.JdbcSinkConnector',
  'tasks.max' = '1',
  'topics' = 'fournewEncode_',
  'connection.url' = 'jdbc:sqlserver://host:port;user={user};password={pass}',
  'connection.user' = {user},
  'connection.password' = {pass},
  'insert.mode' = 'insert',
  'table.name.format' = {table},
  'pk.mode' = 'record_key',
  'pk.fields' = 'user_id',
  'key.converter' = 'org.apache.kafka.connect.json.JsonConverter',
  'key.converter.schemas.enable' = 'false',
  'value.converter'='org.apache.kafka.connect.json.JsonConverter',
  'value.converter.schemas.enable' = 'false',  
  'delete.enabled'='false' -- 关闭删除同步功能
);

方案2:保留删除功能,改用从Value中提取主键

由于Topic的Value中已包含USER_ID字段,可将pk.mode改为record_value,让Connector从Value里读取主键,无需依赖Key的Schema:

CREATE SINK CONNECTOR hsr_user_sqlsrvr1
WITH (
  'connector.class' = 'io.confluent.connect.jdbc.JdbcSinkConnector',
  'tasks.max' = '1',
  'topics' = 'fournewEncode_',
  'connection.url' = 'jdbc:sqlserver://host:port;user={user};password={pass}',
  'connection.user' = {user},
  'connection.password' = {pass},
  'insert.mode' = 'insert',
  'table.name.format' = {table},
  'pk.mode' = 'record_value', -- 切换为从Value提取主键
  'pk.fields' = 'user_id',
  'key.converter' = 'org.apache.kafka.connect.json.JsonConverter',
  'key.converter.schemas.enable' = 'false',
  'value.converter'='org.apache.kafka.connect.json.JsonConverter',
  'value.converter.schemas.enable' = 'false',  
  'delete.enabled'='true'
);

方案3:修改Key的Schema(适合必须保留pk.mode=record_key的场景)

如果一定要用Key作为主键,需要调整上游逻辑让Key携带Schema信息:

  1. 修改源Oracle Connector配置,让它生成带Schema的Struct类型Key(包含user_id字段)
  2. 将Sink Connector的key.converter.schemas.enable改为true,确保Connector能解析Key的Schema结构

注意事项

  • 方案2需保证Topic的Value中USER_ID字段始终非空,否则会导致主键识别失败
  • 若使用insert.mode=insert,重复主键会触发插入失败,可根据业务需求改为upsert实现更新/插入逻辑

内容的提问来源于stack exchange,提问作者Pradyumn Joshi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 21:23:16