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

如何通过Kafka JMS Source Connector将ActiveMQ Artemis队列消息并行写入多个Kafka主题?

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 23:08:13