将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
相关产品推荐
相关产品推荐

