Debezium MySQL源连接器Key侧Subject策略Schema冲突咨询
问题根因
Key侧subject生成不符合预期是两个配置错误叠加导致的:
- Key转换器的subject策略配置项写错了。你配置的
CONNECT_KEY_CONVERTER_KEY_SUBJECT_NAME_STRATEGY是无效配置——不管是Key还是Value的Avro转换器,读取subject命名策略的配置后缀都是value.subject.name.strategy,不存在key.subject.name.strategy这个配置项。配置不生效的情况下,Key转换器默认回退使用TopicNameStrategy,只会按主题名生成单个Key subject,就是你看到的db1_schema.db1_schema-Key。 ByLogicalTableRouter单消息转换(SMT)默认会强制修改Key的Schema名为路由后的新主题名,哪怕策略加载正确,也拿不到原始的服务名.库名.表名格式的记录名,没法按表维度拆分subject。
修复步骤
修正docker-compose环境变量
删除错误的Key侧策略配置项,替换为正确配置,修改后相关配置如下:
- CONNECT_KEY_CONVERTER=io.confluent.connect.avro.AvroConverter - CONNECT_KEY_CONVERTER_SCHEMA_REGISTRY_URL=http://registry:8081 # 注意:Key转换器的subject策略配置后缀和Value转换器一致,均为value.subject.name.strategy - CONNECT_KEY_CONVERTER_VALUE_SUBJECT_NAME_STRATEGY=io.confluent.kafka.serializers.subject.TopicRecordNameStrategy - CONNECT_VALUE_CONVERTER=io.confluent.connect.avro.AvroConverter - CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL=http://registry:8081 - CONNECT_VALUE_CONVERTER_VALUE_SUBJECT_NAME_STRATEGY=io.confluent.kafka.serializers.subject.TopicRecordNameStrategy
补全ByLogicalTableRouter配置参数
在Reroute转换配置中新增key.enforce.uniqueness参数并设为false,禁止SMT覆盖Key的原始Schema名称,修改后的transform相关配置段如下:
"transforms": "unwrap,Reroute", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.delete.handling.mode": "rewrite", "transforms.unwrap.add.fields": "db,table,op,source.ts_ms", "transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter", "transforms.Reroute.topic.regex": "(.*\\S)\\.(.*\\S)\\.(.*\\S)", "transforms.Reroute.topic.replacement": "$2_schema", "transforms.Reroute.key.field.name": "table", "transforms.Reroute.key.field.regex": "(.*\\S)\\.(.*\\S)\\.(.*\\S)", "transforms.Reroute.key.field.replacement": "$3", "transforms.Reroute.key.enforce.uniqueness": "false"
效果验证
配置修改并重启连接器后,Schema Registry会按预期生成独立的Key subject:
- db1_schema.aws-db.db1.table1-Key
- db1_schema.aws-db.db1.table2-Key
- db1_schema.aws-db.db1.table3-Key
哪怕不同表的主键字段类型不一致,也不会触发409 Schema不兼容错误,同一数据库下的所有表消息可以正常路由到对应Kafka主题,Key和Value的Schema都会按表维度独立注册到Schema Registry。
内容的提问来源于stack exchange,提问作者Vektor88
相关产品推荐
相关产品推荐

