【问题标题】:spring boot app to integrate kafka with active mqSpring Boot 应用程序将 kafka 与活动 mq 集成
【发布时间】:2020-04-07 21:32:47
【问题描述】:

我正在尝试构建一个从 kafka 读取消息并将它们放入 activeMQ 的 Spring Boot 应用程序 反之亦然(从activeMQ读取并写入kafka) 我没有找到任何有用的教程来启动我的项目

【问题讨论】:

    标签: spring-boot apache-kafka jms activemq spring-kafka


    【解决方案1】:

    参见Spring IntegrationSpring Integration Extension for Apache Kafka

    使用入站和出站通道适配器

    jms -> kafka
    
    kafka -> jms
    

    Kafka Connect 在这方面也有一些能力,但我不太熟悉。

    编辑

    这个简单的 Spring Boot 应用展示了将数据从 Kafka 传输到 RabbitMQ,反之亦然:

    package com.example.demo;
    
    import org.apache.kafka.clients.admin.NewTopic;
    
    import org.springframework.amqp.core.Queue;
    import org.springframework.amqp.core.QueueBuilder;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.amqp.rabbit.core.RabbitTemplate;
    import org.springframework.beans.factory.annotation.Autowired;
    import org.springframework.boot.ApplicationRunner;
    import org.springframework.boot.SpringApplication;
    import org.springframework.boot.autoconfigure.SpringBootApplication;
    import org.springframework.context.annotation.Bean;
    import org.springframework.kafka.annotation.KafkaListener;
    import org.springframework.kafka.config.TopicBuilder;
    import org.springframework.kafka.core.KafkaTemplate;
    
    @SpringBootApplication
    public class So61069735Application {
    
        public static void main(String[] args) {
            SpringApplication.run(So61069735Application.class, args);
        }
    
        @Autowired
        private KafkaTemplate<String, String> kafkaTemplate;
    
        @Autowired
        private RabbitTemplate rabbitTemplate;
    
        @Bean
        public ApplicationRunner toKafka() {
            return args -> this.kafkaTemplate.send("so61069735-1", "foo");
        }
    
        @KafkaListener(id = "so61069735-1", topics = "so61069735-1")
        public void listen1(String in) {
            System.out.println("From Kafka: " + in);
            this.rabbitTemplate.convertAndSend("so61069735-2", in.toUpperCase());
        }
    
        @RabbitListener(queues = "so61069735-2")
        public void listen2(String in) {
            System.out.println("From Rabbit: " + in);
            this.kafkaTemplate.send("so61069735-3", in + in);
        }
    
        @KafkaListener(id = "so61069735-3", topics = "so61069735-3")
        public void listen(String in) {
            System.out.println("Final: " + in);
        }
    
        @Bean
        public NewTopic topic1() {
            return TopicBuilder.name("so61069735-1").partitions(1).replicas(1).build();
        }
    
        @Bean
        public Queue queue() {
            return QueueBuilder.durable("so61069735-2").build();
        }
    
        @Bean
        public NewTopic topic2() {
            return TopicBuilder.name("so61069735-3").partitions(1).replicas(1).build();
        }
    
    }
    
    spring.kafka.consumer.auto-offset-reset=earliest
    

    结果

    From Kafka: foo
    From Rabbit: FOO
    Final: FOOFOO
    

    【讨论】:

    • 还有其他方法吗? (不使用spring集成)
    • 您可以直接使用 spring-kafka 和 spring-jms 的组合来完成,正如我所说,Kafka Connect 是另一种选择。
    • 你能否更准确地了解涉及 spring-kafka 和 spring-jms 组合的解决方案。你能给我发例子吗
    • 我加了一个例子。
    猜你喜欢
    • 1970-01-01
    • 2019-09-25
    • 2018-10-19
    • 2022-11-01
    • 2019-10-05
    • 2018-09-24
    • 2020-09-18
    • 1970-01-01
    • 2018-01-26
    相关资源
    最近更新 更多