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

如何使用Debezium源连接器配置Kafka主题事件键及解决Schema不匹配问题

核心问题定位

你的字段重命名SMT不生效,本质是三个配置错误:

  1. 字段匹配规则写了不存在的前缀:Debezium输出的业务字段没有public.table1.这类表名前缀,带前缀的匹配规则完全命中不了目标字段。
  2. 没有处理Debezium的原生事件结构:Debezium默认发的消息是嵌套结构,业务字段全放在after/before块里,外层包了op/source/ts_ms等元数据字段,你直接用ReplaceField$Value只会改外层的元数据字段,根本碰不到业务列。
  3. 事件键没有做适配:你没有显式指定键的生成规则,默认生成的键结构和目标主题绑定的Schema不匹配,也没有对键字段做对应的重命名处理。

修复方案

按顺序调整配置即可:

  1. 最前面加ExtractNewRecordState SMT,把after块里的业务数据提取出来作为消息value,去掉Debezium自带的冗余元数据结构,让后续重命名SMT能直接命中业务字段。
  2. 删掉重命名规则里所有public.table1.前缀,直接用数据库原生列名做匹配源。
  3. 新增Key层面的字段重命名SMT,保证消息键的字段名和目标主题的Key Schema一致。
  4. 显式配置message.key.columns参数,指定表的主键列作为消息键的生成依据。
  5. 调整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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 09:21:34