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

如何通过SMT实现Kafka Sink Connector按字段值从单Topic路由至多表

现有配置的问题

你的当前SMT配置存在两个核心问题:

  1. 丢失业务数据:连续使用ExtractField$Value会逐步将整个消息体替换为提取的字段值,最终消息内容只剩TableName字符串,完全丢失了payload.after中需要写入表的业务数据。
  2. 未实现主题路由:RegexRouter的配置regex=.*和replacement=$0等于完全保留原主题名称,没有根据TableName修改目标主题,无法实现路由到不同表的需求。
正确实现方案

要实现从单个Topic路由到多表且自动创建表的需求,需保留原始业务数据的同时,根据payload.source.table的值修改目标主题(让Sink以该主题名作为表名创建表),具体配置如下:

方案1:使用Confluent TemplateRouter SMT(推荐)

TemplateRouter支持直接引用消息中的字段生成新主题名,无需修改原始消息内容:

transforms=routeToTable
transforms.routeToTable.type=org.apache.kafka.connect.transforms.TemplateRouter
# 引用嵌套字段作为新主题名
transforms.routeToTable.template=${payload.source.table}
# 保留原始主题的格式(可选,若不需要可省略)
transforms.routeToTable.topic.format=${topic}

方案2:使用原生SMT组合(无需额外依赖)

如果无法使用Confluent的扩展SMT,可通过InsertField提取表名到消息头部,再结合RegexRouter(注:此方法需要确保原主题名有固定格式,这里假设原主题名为original-topic):

transforms=insertTableField,routeTopic
# 第一步:将source.table的值插入到消息的头部字段
transforms.insertTableField.type=org.apache.kafka.connect.transforms.InsertField$Header
transforms.insertTableField.field=target_table
transforms.insertTableField.value=${payload.source.table}

# 第二步:根据头部字段修改主题名
transforms.routeTopic.type=org.apache.kafka.connect.transforms.RegexRouter
transforms.routeTopic.regex=original-topic
transforms.routeTopic.replacement=${header.target_table}

配套Sink配置(以JDBC Sink为例)

无论使用哪种方案,都需要在Sink连接器中开启自动建表:

auto.create=true
# 若需要自动同步表结构,开启以下配置
auto.evolve=true
# 表名直接使用主题名(因为我们已经将主题替换为TableName)
table.name.format=${topic}

注意事项

  • 确保Kafka Connect集群有权限创建新主题(若使用方案1/2修改主题)。
  • 若使用Debezium CDC源连接器,可直接在Debezium配置中添加transforms,无需额外处理消息格式。
  • 需保证payload.source.table的值符合目标数据库的表名规范(如避免特殊字符)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 14:51:02