如何为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
相关产品推荐
相关产品推荐

