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

Kafka源连接器消息路由至多Topic及相关技术问题咨询

Kafka连接器消息路由与Topic复制方案

一、Source Connector实现消息多Topic路由

无需依赖Streams应用,有两种实用方案:

1. 使用ExtractTopic转换器

这是最轻量化的方案,直接通过内置转换器从消息内容里提取目标Topic名。可选择从消息的Key或Value字段中提取,配置示例如下:

# 启用转换器
transforms=routeToTopic
# 指定从Value中提取Topic字段
transforms.routeToTopic.type=org.apache.kafka.connect.transforms.ExtractTopic$Value
# 消息Value中存储目标Topic名的字段(示例为target_topic)
transforms.routeToTopic.field=target_topic

若提取出的Topic不存在,只要Kafka集群开启了auto.create.topics.enable配置,Connector会自动创建该Topic。

2. 自定义Source逻辑

如果需要更复杂的路由规则(比如基于多个字段组合判断),可在自定义Source Connector的代码中,直接为每条消息指定目标Topic,无需依赖转换器。

二、复制/移动Topic A到Topic B的Sink方案

有两种现成的Connector可实现需求,无需额外开发:

1. Kafka Sink Connector(同集群复制)

Kafka官方提供的KafkaSinkConnector可直接将一个Topic的消息转发到另一个Topic,配置简单:

name=topic-copy-sink
connector.class=org.apache.kafka.connect.sink.KafkaSinkConnector
tasks.max=1
# 源Topic
topics=TopicA
# 固定目标Topic
sink.topic=TopicB

若需要动态转发到不同Topic,可结合ExtractTopic转换器(参考上方配置),实现基于消息内容的动态路由。

2. MirrorSourceConnector(跨/同集群复制)

MirrorMaker 2.0的MirrorSourceConnector不仅支持跨集群Topic复制,也适合同集群内的Topic复制,适合需要同步Topic元数据(如分区数、副本数)的场景:

name=mirror-topic-connector
connector.class=org.apache.kafka.connect.mirror.MirrorSourceConnector
tasks.max=1
source.cluster.alias=local
target.cluster.alias=local
# 要复制的源Topic
topics=TopicA
# 源Topic到目标Topic的映射
topic.map=TopicA:TopicB

关于“移动”消息:没有现成Connector直接支持“移动”(复制后删除源消息),可通过给源Topic设置短retention.ms自动清理旧消息,或复制完成后手动清理源Topic数据。

三、ExtractTopic转换器的能力

ExtractTopic完全可以满足你的需求:

  • 在Source端:从源数据中提取字段作为发送的Topic,实现动态路由到多个Topic。
  • 在Sink端:配合Kafka Sink Connector使用,从消息中提取字段作为目标Topic,实现从Topic A到不同目标Topic的分发(包括Topic B)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 17:32:30