如何使用Debezium源连接器配置Kafka主题事件键及解决Schema不匹配问题
核心问题定位
你的字段重命名SMT不生效,本质是三个配置错误:
- 字段匹配规则写了不存在的前缀:Debezium输出的业务字段没有
public.table1.这类表名前缀,带前缀的匹配规则完全命中不了目标字段。 - 没有处理Debezium的原生事件结构:Debezium默认发的消息是嵌套结构,业务字段全放在
after/before块里,外层包了op/source/ts_ms等元数据字段,你直接用ReplaceField$Value只会改外层的元数据字段,根本碰不到业务列。 - 事件键没有做适配:你没有显式指定键的生成规则,默认生成的键结构和目标主题绑定的Schema不匹配,也没有对键字段做对应的重命名处理。
修复方案
按顺序调整配置即可:
- 最前面加
ExtractNewRecordStateSMT,把after块里的业务数据提取出来作为消息value,去掉Debezium自带的冗余元数据结构,让后续重命名SMT能直接命中业务字段。 - 删掉重命名规则里所有
public.table1.前缀,直接用数据库原生列名做匹配源。 - 新增Key层面的字段重命名SMT,保证消息键的字段名和目标主题的Key Schema一致。
- 显式配置
message.key.columns参数,指定表的主键列作为消息键的生成依据。 - 调整SMT执行顺序:先提取扁平化记录,再分别重命名Value和Key的字段,最后执行主题路由。
修正后可直接用的配置
{ "name": "source-1", "connector.class": "io.debezium.connector.postgresql.PostgresConnector", "database.hostname": "********************", "database.user": "******************", "database.password": "*************", "database.dbname": "*************", "database.server.name": "**********", "plugin.name": "pgoutput", "table.include.list": "public.table1", "column.include.list": "public.table1.transferId, public.table1.originSpaceId, public.table1.originalManuscriptId, public.table1.lastModifiedDate, public.table1.finalManuscriptId, public.table1.destinationSpaceId", "slot.name": "graphv1_table11", "publication.autocreate.mode": "filtered", "snapshot.mode": "initial", "tasks.max": "1", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.basic.auth.credentials.source": "USER_INFO", "value.converter.basic.auth.user.info": "*********:**************", "value.converter.schema.registry.url": "**********", "key.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.basic.auth.credentials.source": "USER_INFO", "key.converter.basic.auth.user.info": "********:*******************", "key.converter.schema.registry.url": "****************", "message.key.columns": "public.table1:transferId", "transforms": "ExtractRecord,RenameValueField,RenameKeyField,Reroute", "transforms.ExtractRecord.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.ExtractRecord.drop.tombstones": "false", "transforms.ExtractRecord.delete.handling.mode": "rewrite", "transforms.RenameValueField.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.RenameValueField.renames": "transferId:id,originSpaceId:originalSpaceId,originalManuscriptId:originalArticleId,lastModifiedDate:modificationDate,destinationSpaceId:finalSpaceId,finalManuscriptId:finalArticleId", "transforms.RenameKeyField.type": "org.apache.kafka.connect.transforms.ReplaceField$Key", "transforms.RenameKeyField.renames": "transferId:id", "transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter", "transforms.Reroute.topic.regex": ".+", "transforms.Reroute.topic.replacement": "mydestinationTopic" }
注意事项
如果目标主题已经在Schema Registry注册了Avro Schema,必须保证重命名后的字段类型、字段顺序、空值约束和已有Schema完全兼容,否则Avro序列化阶段会直接报错。如果之前已经写入过不兼容的Schema,可以先调整Schema Registry的兼容策略,或者删除主题对应的错误Schema版本再启动任务。
如果你的主键不是
transferId,把message.key.columns和RenameKeyField.renames里的对应字段改成实际主键即可。
内容的提问来源于stack exchange,提问作者Muhammed Rahiz
相关产品推荐
相关产品推荐

