如何使用Spring Boot或Spring Kafka实现多源主题到多目标主题的消息传输
Spring Kafka多源到多目标主题消息传输替代实现方案
除手动遍历目标主题逐一发送外,还有三种常见的低代码、高性能实现方案:
方案1:Spring Integration Kafka 声明式配置
无需手动编写消费、发送的核心业务逻辑,仅通过配置即可完成消息流转,适合固定规则的转发场景:
- 首先引入依赖:
spring-boot-starter-integration、spring-integration-kafka - 配置多源主题入站适配器,一次性监听所有源主题:
@Bean public KafkaMessageDrivenChannelAdapter<String, String> kafkaInboundAdapter( ConsumerFactory<String, String> consumerFactory) { ContainerProperties containerProps = new ContainerProperties("source-topic-1", "source-topic-2", "source-topic-3"); KafkaMessageDrivenChannelAdapter<String, String> adapter = new KafkaMessageDrivenChannelAdapter<>(new ConcurrentMessageListenerContainer<>(consumerFactory, containerProps)); adapter.setOutputChannel(kafkaSendChannel()); return adapter; }
- 配置多目标出站处理,新增/删除目标主题仅需增减对应Bean即可:
// 发布订阅通道会将收到的消息转发给所有绑定的处理器 @Bean public PublishSubscribeChannel kafkaSendChannel() { return new PublishSubscribeChannel(); } // 目标主题1出站处理器 @Bean @ServiceActivator(inputChannel = "kafkaSendChannel") public KafkaProducerMessageHandler<String, String> target1Handler( KafkaTemplate<String, String> kafkaTemplate) { KafkaProducerMessageHandler<String, String> handler = new KafkaProducerMessageHandler<>(kafkaTemplate); handler.setTopicExpression(new LiteralExpression("target-topic-1")); return handler; } // 目标主题2出站处理器 @Bean @ServiceActivator(inputChannel = "kafkaSendChannel") public KafkaProducerMessageHandler<String, String> target2Handler( KafkaTemplate<String, String> kafkaTemplate) { KafkaProducerMessageHandler<String, String> handler = new KafkaProducerMessageHandler<>(kafkaTemplate); handler.setTopicExpression(new LiteralExpression("target-topic-2")); return handler; }
方案2:Kafka Streams API 集成
适合需要对消息做过滤、转换后再转发的大流量场景,Spring Kafka原生支持Kafka Streams生命周期管理,性能远高于手动编写的转发逻辑:
- 引入依赖:
spring-boot-starter-kafka-streams - 配置流转发规则:
@Configuration @EnableKafkaStreams public class KafkaStreamConfig { @Bean public KStream<String, String> forwardStream(StreamsBuilder streamsBuilder) { // 同时监听多个源主题 KStream<String, String> stream = streamsBuilder.stream(List.of("source-topic-1", "source-topic-2")); // 可在此处添加filter、map等消息处理逻辑 // 批量转发到多个目标主题 stream.to("target-topic-1"); stream.to("target-topic-2"); stream.to("target-topic-3"); return stream; } }
方案3:KafkaTemplate 异步批量发送优化
如果你希望保留代码自定义能力,可将逐次同步发送改为批量异步提交,降低线程阻塞开销:
@Autowired private KafkaTemplate<String, String> kafkaTemplate; // 监听多个源主题 @KafkaListener(topics = {"source-topic-1", "source-topic-2"}) public void onMessage(ConsumerRecord<String, String> record) { List<String> targetTopics = List.of("target-topic-1", "target-topic-2", "target-topic-3"); // 批量提交所有发送请求,无需等待单条发送完成 List<CompletableFuture<SendResult<String, String>>> sendFutures = targetTopics.stream() .map(topic -> kafkaTemplate.send(topic, record.key(), record.value())) .toList(); // 可选:统一监听所有发送结果做异常处理 CompletableFuture.allOf(sendFutures.toArray(new CompletableFuture[0])) .whenComplete((res, e) -> { if (e != null) { // 自定义发送失败处理逻辑 } }); }
内容的提问来源于stack exchange,提问作者Ram
相关产品推荐
相关产品推荐

