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

将Camel中ActiveMQ路由及Bean迁移至Kafka的正确方法

Camel ActiveMQ路由迁移至Kafka的正确方式

现有ActiveMQ配置

路由定义

<route id="readInactiveRecordNotification">
        <from uri="appActiveMQConsumer:queue:Consumer.hydra.ldap_updates_poller.VirtualTopic.ldap.in" />
</route>

对应的ActiveMQ组件Bean

@Bean(name = "appActiveMQConsumer")
public ActiveMQComponent amqConsumer() throws Exception {
    ActiveMQComponent amqConsumer = new ActiveMQComponent();
    amqConsumer.setConnectionFactory(pooledConnectionFactoryConsumer());
    amqConsumer.setTransacted(true);
    amqConsumer.setDeliveryPersistent(true);
    amqConsumer.setCacheLevelName("CACHE_CONSUMER");
    amqConsumer.setAcknowledgementModeName("SESSION_TRANSACTED");
    return amqConsumer;
}

已准备的Kafka组件Bean

@Bean(name="appkafkaConsumer")
public KafkaComponent kafkaComponentConsumer(){
    KafkaComponent kafkaComponentConsumer = new KafkaComponent();
    kafkaComponentConsumer.setConfiguration(kafkaConfiguration());
    return kafkaComponentConsumer;
}

迁移后的Kafka路由配置

XML路由实现

<route id="readInactiveRecordNotification">
    <from uri="appkafkaConsumer:ldap.in?groupId=hydra.ldap_updates_poller" />
</route>

Java DSL路由实现

from("appkafkaConsumer:ldap.in?groupId=hydra.ldap_updates_poller")
    .routeId("readInactiveRecordNotification");

核心迁移要点

  • ActiveMQ的VirtualTopic.ldap.in对应Kafka的主题ldap.in
  • ActiveMQ队列名称中的Consumer.hydra.ldap_updates_poller对应Kafka的消费者组hydra.ldap_updates_poller,Kafka通过消费者组实现类似VirtualTopic的多订阅者模式
  • 事务一致性:原ActiveMQ开启了会话事务,需确保Kafka配置中关闭自动提交(enable.auto.commit=false),并启用事务型消费者,和原有逻辑保持一致
  • 消息持久化:Kafka默认开启消息持久化,无需额外配置即可对应原ActiveMQ的deliveryPersistent=true
  • 缓存机制:Kafka消费者本身维护连接与订阅缓存,原ActiveMQ的CACHE_CONSUMER无需额外配置

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 11:45:37