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

Kafka Connect ExtractField$Key SMT报Unknown Field错误问题咨询

故障原因

报错的核心原因有两点:

  • Debezium SQL Server连接器默认开启include.schema.changes配置(默认值为true),同步任务启动时会先向输出Topic发送Schema变更元数据事件。这类事件的消息Key仅包含数据库名、Schema名、表名字段,不存在你要提取的InternalID字段,ExtractField SMT检测到目标字段不存在时直接抛出异常。
  • 你当前配置的SMT执行顺序为unwrap -> extractField,其中ExtractNewRecordState(即unwrap SMT)仅会对增删改类型的业务CRUD事件做结构展开,对Schema变更这类元数据事件会直接透传,不做任何结构转换,导致后续的extractField直接处理不符合预期结构的元数据消息,触发报错。

移除extractField配置后看起来Key正常,是因为元数据事件仅在任务启动、表结构变更时少量发送,你消费时大概率只看到了后续的业务数据消息,没有感知到这类特殊结构的消息。

解决方案

按以下步骤调整配置即可:

  1. 新增配置项关闭Schema变更元数据事件输出,避免非业务消息进入SMT处理链
  2. (推荐)调整SMT执行顺序,将Key提取操作放到unwrap之前,避免SMT版本逻辑变更导致Key结构被意外修改

修改后的完整连接器配置如下:

CREATE SOURCE CONNECTOR properties_sql_connector WITH (
'connector.class'= 'io.debezium.connector.sqlserver.SqlServerConnector', 
'database.hostname'= 'propertiessql', 
'database.port'= '1433', 
'database.user'= 'XXX', 
'database.password'= 'XXX', 
'database.dbname'= 'Properties', 
'database.server.name'= 'properties', 
'table.exclude.list'= 'dbo.__EFMigrationsHistory', 
'database.history.kafka.bootstrap.servers'= 'kafka:9091', 
'database.history.kafka.topic'= 'dbhistory.properties',
'key.converter.schemas.enable'= 'false',
'include.schema.changes'= 'false',
'transforms'= 'extractField,unwrap',
'transforms.unwrap.type'= 'io.debezium.transforms.ExtractNewRecordState',
'transforms.unwrap.delete.handling.mode'= 'none',
'transforms.extractField.type'= 'org.apache.kafka.connect.transforms.ExtractField$Key',
'transforms.extractField.field'= 'InternalID',
'key.converter'= 'org.apache.kafka.connect.json.JsonConverter'
);

如果需要保留Schema变更事件,可以额外配置Kafka Connect Predicate规则,让extractField仅对业务CRUD事件生效,普通CDC同步场景直接关闭Schema变更事件即可满足需求,不会影响业务数据的同步准确性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 20:27:25