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正常,是因为元数据事件仅在任务启动、表结构变更时少量发送,你消费时大概率只看到了后续的业务数据消息,没有感知到这类特殊结构的消息。
解决方案
按以下步骤调整配置即可:
- 新增配置项关闭Schema变更元数据事件输出,避免非业务消息进入SMT处理链
- (推荐)调整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
相关产品推荐
相关产品推荐

