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信息:
- 修改源Oracle Connector配置,让它生成带Schema的Struct类型Key(包含
user_id字段) - 将Sink Connector的
key.converter.schemas.enable改为true,确保Connector能解析Key的Schema结构
注意事项
- 方案2需保证Topic的Value中
USER_ID字段始终非空,否则会导致主键识别失败 - 若使用
insert.mode=insert,重复主键会触发插入失败,可根据业务需求改为upsert实现更新/插入逻辑
内容的提问来源于stack exchange,提问作者Pradyumn Joshi
相关产品推荐
相关产品推荐

