【发布时间】:2020-04-07 21:32:47
【问题描述】:
我正在尝试构建一个从 kafka 读取消息并将它们放入 activeMQ 的 Spring Boot 应用程序 反之亦然(从activeMQ读取并写入kafka) 我没有找到任何有用的教程来启动我的项目
【问题讨论】:
标签: spring-boot apache-kafka jms activemq spring-kafka
我正在尝试构建一个从 kafka 读取消息并将它们放入 activeMQ 的 Spring Boot 应用程序 反之亦然(从activeMQ读取并写入kafka) 我没有找到任何有用的教程来启动我的项目
【问题讨论】:
标签: spring-boot apache-kafka jms activemq spring-kafka
参见Spring Integration 和Spring 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
【讨论】: