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

如何为Debezium PostgreSQL连接器的所有主题配置统一路由规则

问题场景

我的PostgreSQL数据库包含以下表:

mdm.table1_a 
mdm.table2_a_b 
mdm.table3_a
mdm.table4_a_b_c
mdm.table5.a_b
mdm.table6_a_b_c

Debezium PostgreSQL连接器会结合topic.prefix(示例值为"ips")生成与表名对应的主题:

ips.mdm.table1_a 
ips.mdm.table2_a_b 
ips.mdm.table3_a
ips.mdm.table4_a_b_c
ips.mdm.table5.a_b
ips.mdm.table6_a_b_c

我需要将主题名中的下划线替换为连字符,目前通过Debezium的ByLogicalTableRouter转换配置了3条规则:

"transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter",
"transforms.Reroute.topic.regex": "(.*\\.mdm\\.)([a-z0-9]*)_([a-z0-9]*)_([a-z0-9]*)",       
"transforms.Reroute.topic.replacement": "$1$2-$3-$4",       
"transforms.Reroute2.type": "io.debezium.transforms.ByLogicalTableRouter",
"transforms.Reroute2.topic.regex": "(.*\\.mdm\\.)([a-z0-9]*)_([a-z0-9]*)",
"transforms.Reroute2.topic.replacement": "$1$2-$3",
"transforms.Reroute3.type": "io.debezium.transforms.ByLogicalTableRouter",
"transforms.Reroute3.topic.regex": "(.*\\.mdm\\.)([a-z0-9]*)_([a-z0-9]*)_([a-z0-9]*)_([a-z0-9]*)",
"transforms.Reroute3.topic.replacement": "$1$2-$3-$4-$5"

现在想知道能不能把这3条规则合并成一个通用的路由规则?

完整的连接器配置如下:

{
  "name": "debezium_mdm",
  "config": {
      "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
      "database.hostname": "pg-wilddev.site1.com",
      "database.port": "5432",
      "database.dbname": "mdm",
      "database.user": "debezium",
      "database.password": "",
      "schema.include.list": "mdm",   
      "table.include.list": "mdm.table1_a,mdm.table2_a_b,mdm.table3_a,mdm.table4_a_b_c,mdm.table5.a_b,mdm.table6_a_b_c",
      "topic.prefix": "ips",
      "slot.name": "dbz_mdm",
      "publication.name": "dbz_mdm",
      "plugin.name": "pgoutput",
      "time.precision.mode": "connect",
      "heartbeat.interval.ms": "5000",
      "retriable.restart.connector.wait.ms": "100000",
      "topic.creation.default.partitions": "3",
      "topic.creation.default.replication.factor": "3",
      "topic.creation.groups": "mdm", 
      "topic.creation.mdm.partitions": "2",       
      "topic.creation.mdm.replication.factor": "3",
      "topic.creation.mdm.compression.type": "lz4",
      "topic.creation.mdm.include": ".*mdm.*",
      "topic.creation.mdm.cleanup.policy": "compact",
      "key.converter": "org.apache.kafka.connect.storage.StringConverter",
      "key.converter.schemas.enable": "false",
      "transforms": "unwrap, extractToKey, Reroute, Reroute2, Reroute3",
      "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
      "transforms.unwrap.delete.tombstone.handling.mode": "drop",
      "transforms.extractToKey.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
      "transforms.extractToKey.field": "id",      
      "predicates": "topicNameMatch",
      "predicates.topicNameMatch.type": "org.apache.kafka.connect.transforms.predicates.TopicNameMatches",
      "predicates.topicNameMatch.pattern": "ips.mdm.*",
      "transforms.extractToKey.predicate": "topicNameMatch",
      "transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter",
      "transforms.Reroute.topic.regex": "(.*\\.mdm\\.)([a-z0-9]*)_([a-z0-9]*)_([a-z0-9]*)",       
      "transforms.Reroute.topic.replacement": "$1$2-$3-$4",       
      "transforms.Reroute2.type": "io.debezium.transforms.ByLogicalTableRouter",
      "transforms.Reroute2.topic.regex": "(.*\\.mdm\\.)([a-z0-9]*)_([a-z0-9]*)",
      "transforms.Reroute2.topic.replacement": "$1$2-$3",
      "transforms.Reroute3.type": "io.debezium.transforms.ByLogicalTableRouter",
      "transforms.Reroute3.topic.regex": "(.*\\.mdm\\.)([a-z0-9]*)_([a-z0-9]*)_([a-z0-9]*)_([a-z0-9]*)",
      "transforms.Reroute3.topic.replacement": "$1$2-$3-$4-$5"
      }
}
解决方案

可以通过单一正则表达式规则实现所有下划线替换为连字符的需求,无需拆分多个规则。核心思路是匹配mdm.之后的所有内容,将其中的下划线全局替换为连字符。

修改后的转换配置

"transforms": "unwrap, extractToKey, Reroute",
"transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter",
"transforms.Reroute.topic.regex": "(.*\\.mdm\\.)(.*)",
"transforms.Reroute.topic.replacement": "$1${2/_/-g}"

配置说明

  • (.*\\.mdm\\.):匹配主题名中mdm.之前的所有部分(包含mdm.),捕获为分组1。
  • (.*):匹配mdm.之后的表名部分,捕获为分组2。
  • ${2/_/-g}:对分组2的内容执行全局替换,将所有下划线_替换为连字符-。g表示全局匹配,确保替换所有符合条件的字符。

替换后的完整连接器配置

{
  "name": "debezium_mdm",
  "config": {
      "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
      "database.hostname": "pg-wilddev.site1.com",
      "database.port": "5432",
      "database.dbname": "mdm",
      "database.user": "debezium",
      "database.password": "",
      "schema.include.list": "mdm",   
      "table.include.list": "mdm.table1_a,mdm.table2_a_b,mdm.table3_a,mdm.table4_a_b_c,mdm.table5.a_b,mdm.table6_a_b_c",
      "topic.prefix": "ips",
      "slot.name": "dbz_mdm",
      "publication.name": "dbz_mdm",
      "plugin.name": "pgoutput",
      "time.precision.mode": "connect",
      "heartbeat.interval.ms": "5000",
      "retriable.restart.connector.wait.ms": "100000",
      "topic.creation.default.partitions": "3",
      "topic.creation.default.replication.factor": "3",
      "topic.creation.groups": "mdm", 
      "topic.creation.mdm.partitions": "2",       
      "topic.creation.mdm.replication.factor": "3",
      "topic.creation.mdm.compression.type": "lz4",
      "topic.creation.mdm.include": ".*mdm.*",
      "topic.creation.mdm.cleanup.policy": "compact",
      "key.converter": "org.apache.kafka.connect.storage.StringConverter",
      "key.converter.schemas.enable": "false",
      "transforms": "unwrap, extractToKey, Reroute",
      "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
      "transforms.unwrap.delete.tombstone.handling.mode": "drop",
      "transforms.extractToKey.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
      "transforms.extractToKey.field": "id",      
      "predicates": "topicNameMatch",
      "predicates.topicNameMatch.type": "org.apache.kafka.connect.transforms.predicates.TopicNameMatches",
      "predicates.topicNameMatch.pattern": "ips.mdm.*",
      "transforms.extractToKey.predicate": "topicNameMatch",
      "transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter",
      "transforms.Reroute.topic.regex": "(.*\\.mdm\\.)(.*)",
      "transforms.Reroute.topic.replacement": "$1${2/_/-g}"
      }
}

效果验证

应用该配置后,原主题名将被转换为:

ips.mdm.table1-a 
ips.mdm.table2-a-b 
ips.mdm.table3-a
ips.mdm.table4-a-b-c
ips.mdm.table5.a-b
ips.mdm.table6-a-b-c

完全符合将所有下划线替换为连字符的需求,同时保留了原表名中的点号(如table5.a_b转换为table5.a-b)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 03:57:06