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

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组合,通过独立连接器分别处理不同事件类型:

  1. 创建第一个连接器,专门过滤插入事件并路由到对应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
  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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 13:24:55