Debezium能否将单表记录路由至多个Kafka Topic?
解决同表不同事件路由到不同Kafka Topic的问题
核心问题
多个ByLogicalTableRouter针对同一张表配置时,会触发Debezium规则冲突——Debezium无法同时处理多个针对同一表的路由逻辑,直接导致任务启动失败或消息处理异常。
正确实现方式
方法1:自定义单一路由器,按事件类型分流
编写一个自定义ByLogicalTableRouter实现,在路由逻辑里根据Debezium事件的op字段(操作类型标识)判断目标Topic:
public class EventTypeBasedRouter extends ByLogicalTableRouter { @Override public String route(String record, DatabaseSchema schema, TableId tableId) { // 解析Debezium变更事件,提取操作类型 JsonNode json = JsonUtil.parse(record); String op = json.get("payload").get("op").asText(); // 根据操作类型路由到对应Topic if ("c".equals(op)) { // c = 插入(创建)操作 return "topic-insert-" + tableId.table(); } else if ("d".equals(op)) { // d = 删除操作 return "topic-delete-" + tableId.table(); } // 更新(u)等其他操作可自定义默认Topic,或沿用父类逻辑 return super.route(record, schema, tableId); } }
然后在Debezium配置中仅配置这一个路由器:
transforms=route transforms.route.type=com.yourpackage.EventTypeBasedRouter transforms.route.topic.regex=(.*) transforms.route.topic.replacement=$1
方法2:用内置SMT结合多连接器分流
利用Debezium的Filter和TopicRouter组合,通过独立连接器分别处理不同事件类型:
- 创建第一个连接器,专门过滤插入事件并路由到对应Topic:
transforms=filterInsert,routeInsert transforms.filterInsert.type=io.debezium.transforms.Filter transforms.filterInsert.language=jsr223.groovy transforms.filterInsert.condition=payload.op == 'c' transforms.routeInsert.type=io.debezium.transforms.ByLogicalTableRouter transforms.routeInsert.topic.regex=(.*) transforms.routeInsert.topic.replacement=topic-insert-$1
- 创建第二个独立连接器,专门过滤删除事件并路由到对应Topic:
transforms=filterDelete,routeDelete transforms.filterDelete.type=io.debezium.transforms.Filter transforms.filterDelete.language=jsr223.groovy transforms.filterDelete.condition=payload.op == 'd' transforms.routeDelete.type=io.debezium.transforms.ByLogicalTableRouter transforms.routeDelete.topic.regex=(.*) transforms.routeDelete.topic.replacement=topic-delete-$1
注意:此方案需为同一张表配置两个独立连接器,避免路由规则冲突,同时要做好连接器的偏移量管理,防止重复或漏消费。
关键注意事项
- 禁止为同一张表配置多个
ByLogicalTableRouter,Debezium的SMT按顺序执行,重复规则会直接引发逻辑混乱。 - 自定义路由器时要注意不同数据库的Debezium事件格式差异,比如MySQL和PostgreSQL的payload字段结构略有不同,需针对性解析。
- 多连接器方案中,要确保两个连接器的表监听范围一致,避免出现事件遗漏。
内容的提问来源于stack exchange,提问作者Ronn M
相关产品推荐
相关产品推荐

