如何通过Kafka JMS Source Connector将ActiveMQ Artemis队列消息并行写入多个Kafka主题?
咱们先戳破核心问题:JMS队列是点对点模型,消息被一个消费者读取后就会从队列中删除,所以你部署多个Connector实例时,它们会竞争消费,每个消息只会被其中一个实例拿到;而默认的RegexRouter SMT只能修改单条消息的目标主题,没法复制消息发送到多个主题。下面给你三个可行的解决方案,按需选择:
方案1:利用ActiveMQ Artemis原生消息复制功能(推荐,无代码改动)
这个思路是先让Artemis把源队列的消息复制到多个独立的队列,再给每个队列配一个Kafka Connector,这样每个Connector都能拿到完整的消息集,各自写入对应的Kafka主题。
步骤1:配置Artemis的Divert(分流)实现消息复制
在Artemis的broker.xml中添加Divert配置,把源队列的消息复制(不是移动)到多个目标队列:
<diverts> <!-- 复制源队列消息到topic1专属队列 --> <divert name="divert-to-topic1-queue"> <address>your-source-queue</address> <!-- 替换成你的源队列地址 --> <forwarding-address>topic1-target-queue</address> <!-- 目标队列1 --> <exclusive>false</exclusive> <!-- 关键:false=复制,true=移动 --> <routing-type>MULTICAST</routing-type> </divert> <!-- 复制源队列消息到topic2专属队列 --> <divert name="divert-to-topic2-queue"> <address>your-source-queue</address> <forwarding-address>topic2-target-queue</address> <!-- 目标队列2 --> <exclusive>false</exclusive> <routing-type>MULTICAST</routing-type> </divert> <!-- 按需添加更多divert,对应更多Kafka主题 --> </diverts>
配置完重启Artemis,源队列的消息就会自动复制到所有目标队列。
步骤2:为每个目标队列配置独立的Kafka Connector
每个Connector指向一个目标队列,写入对应的Kafka主题。比如第一个Connector的配置:
name=artemis-to-kafka-topic1 connector.class=io.confluent.connect.jms.JmsSourceConnector tasks.max=1 kafka.topic=your-kafka-topic-1 jms.url=tcp://your-artemis-host:61616 jms.username=your-username jms.password=your-password jms.destination.name=topic1-target-queue jms.destination.type=QUEUE
第二个Connector只需要修改name、kafka.topic和jms.destination.name即可,其他配置保持一致。
方案2:在Kafka Connect层面实现多主题转发(需自定义扩展)
如果无法修改Artemis配置,可以在Kafka Connect侧实现消息复制,有两种方式:
方式A:自定义ProducerInterceptor
编写一个Kafka Producer拦截器,在消息发送前复制消息并发送到多个主题。示例代码(Java):
public class MultiTopicProducerInterceptor implements ProducerInterceptor<String, byte[]> { private List<String> targetTopics; @Override public void configure(Map<String, ?> configs) { targetTopics = Arrays.asList(configs.get("target.topics").toString().split(",")); } @Override public ProducerRecord<String, byte[]> onSend(ProducerRecord<String, byte[]> record) { // 先发送原消息到默认主题 targetTopics.forEach(topic -> { if (!topic.equals(record.topic())) { ProducerRecord<String, byte[]> copyRecord = new ProducerRecord<>( topic, record.key(), record.value(), record.headers() ); // 用Producer发送复制的消息(需要持有Producer实例,这里简化逻辑) ((KafkaProducer<String, byte[]>) context.producer()).send(copyRecord); } }); return record; } // 实现其他Interceptor方法 @Override public void onAcknowledgement(RecordMetadata metadata, Exception exception) {} @Override public void close() {} }
然后在Connector配置中添加拦截器:
producer.interceptor.classes=com.yourcompany.MultiTopicProducerInterceptor target.topics=kafka-topic-1,kafka-topic-2,kafka-topic-3
方式B:使用社区多主题SMT
有些社区维护的SMT支持多主题转发,你可以在Confluent Hub上搜索multi-topic router,找到现成的组件直接使用,不需要自己写代码。
方案3:先用Connector写入单Kafka主题,再用Kafka Streams复制到多主题(解耦架构)
这个方案适合已经在用Kafka生态的场景,把消息复制的逻辑和JMS消费逻辑解耦:
步骤1:配置单个Connector写入中间主题
先把Artemis队列的消息全部写入一个中间Kafka主题,比如jms-raw-messages:
name=artemis-to-kafka-middle connector.class=io.confluent.connect.jms.JmsSourceConnector tasks.max=1 kafka.topic=jms-raw-messages jms.url=tcp://your-artemis-host:61616 jms.username=your-username jms.password=your-password jms.destination.name=your-source-queue jms.destination.type=QUEUE
步骤2:用Kafka Streams实现多主题复制
编写一个简单的Streams应用,读取中间主题并转发到多个目标主题:
import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.KStream; import java.util.Properties; public class JmsMessageTopicCopier { public static void main(String[] args) { Properties streamsProps = new Properties(); streamsProps.put(StreamsConfig.APPLICATION_ID_CONFIG, "jms-message-copier-app"); streamsProps.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-broker:9092"); StreamsBuilder builder = new StreamsBuilder(); KStream<String, byte[]> sourceStream = builder.stream("jms-raw-messages"); // 转发到多个目标主题 sourceStream.to("kafka-topic-1"); sourceStream.to("kafka-topic-2"); sourceStream.to("kafka-topic-3"); KafkaStreams streams = new KafkaStreams(builder.build(), streamsProps); streams.start(); // 优雅关闭钩子 Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); } }
这个方案的好处是,你还可以在Streams里添加过滤、转换等逻辑,灵活性很高。
内容的提问来源于stack exchange,提问作者topicfun

