基于RabbitMQ的Super Streams与Spring Cloud Stream消费组实现示例求助
基于RabbitMQ实现多服务消费同数据且单服务实例不重复消费的方案
一、RabbitMQ Super Streams 实现示例
Super Streams是RabbitMQ的分片流队列集合,核心逻辑是:
- 单个Super Stream包含多个分片队列,生产者按规则将消息分发到不同分片
- 同一消费者组内的实例会分摊消费不同分片,天然避免组内重复消费
- 不同服务可创建独立消费者组,各自消费全量数据
1. 前置准备
确保RabbitMQ已启用stream插件,执行命令:
rabbitmq-plugins enable rabbitmq_stream
2. 创建Super Stream
通过CLI创建包含3个分片的Super Stream(名称为common-data-stream):
rabbitmqadmin declare super_stream name=common-data-stream partitions=3
3. 生产者代码示例(Java)
使用RabbitMQ Stream客户端发送消息,消息会自动分片:
import com.rabbitmq.stream.Environment; import com.rabbitmq.stream.Producer; import com.rabbitmq.stream.ProducerBuilder; public class SuperStreamProducer { public static void main(String[] args) { Environment environment = Environment.builder().build(); Producer producer = ProducerBuilder .producer(environment) .stream("common-data-stream") .build(); // 模拟发送10条测试数据 for (int i = 0; i < 10; i++) { producer.send(environment.messageBuilder() .addData(("Common Data " + i).getBytes()) .build()); } producer.close(); environment.close(); } }
4. 服务A消费者组(实例分摊消费)
创建消费者组service-a-group,组内多个实例会自动分摊不同分片的消息:
import com.rabbitmq.stream.Environment; import com.rabbitmq.stream.Consumer; import com.rabbitmq.stream.ConsumerBuilder; public class ServiceAConsumer { public static void main(String[] args) { Environment environment = Environment.builder().build(); Consumer consumer = ConsumerBuilder .consumer(environment) .stream("common-data-stream") .consumerGroup("service-a-group") .name("service-a-instance-" + args[0]) // 传入实例标识,如A1、A2 .messageHandler((context, message) -> { System.out.println("Service A Instance " + args[0] + " received: " + new String(message.getBodyAsBinary())); }) .build(); // 注册关闭钩子,优雅停止消费 Runtime.getRuntime().addShutdownHook(new Thread(consumer::close)); } }
5. 服务B消费者组(独立消费全量数据)
创建独立消费者组service-b-group,与服务A组互不干扰,消费全量数据:
import com.rabbitmq.stream.Environment; import com.rabbitmq.stream.Consumer; import com.rabbitmq.stream.ConsumerBuilder; public class ServiceBConsumer { public static void main(String[] args) { Environment environment = Environment.builder().build(); Consumer consumer = ConsumerBuilder .consumer(environment) .stream("common-data-stream") .consumerGroup("service-b-group") .name("service-b-instance-" + args[0]) .messageHandler((context, message) -> { System.out.println("Service B Instance " + args[0] + " received: " + new String(message.getBodyAsBinary())); }) .build(); Runtime.getRuntime().addShutdownHook(new Thread(consumer::close)); } }
二、Spring Cloud Stream + RabbitMQ 实现示例
Spring Cloud Stream通过绑定器和消费者组配置,原生支持需求:
- 同一消费者组内的实例分摊消费消息,避免重复
- 不同服务使用不同消费者组,各自消费全量数据
1. 依赖配置(Maven)
<dependencies> <dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-stream-rabbit</artifactId> </dependency> </dependencies>
2. 服务A配置(application.yml)
指定消费者组service-a-group,确保组内实例分摊消费:
spring: cloud: stream: bindings: input: destination: common-data-exchange # 共享交换器 group: service-a-group # 消费者组 binder: rabbit rabbit: bindings: input: consumer: auto-bind-dlq: true # 可选:自动绑定死信队列
3. 服务A消费者代码
import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.messaging.Sink; import org.springframework.stereotype.Component; @Component public class ServiceAConsumer { @StreamListener(Sink.INPUT) public void handleMessage(String message) { System.out.println("Service A Instance received: " + message); } }
4. 服务B配置(application.yml)
使用独立消费者组service-b-group:
spring: cloud: stream: bindings: input: destination: common-data-exchange group: service-b-group binder: rabbit rabbit: bindings: input: consumer: auto-bind-dlq: true
5. 服务B消费者代码
import org.springframework.cloud.stream.annotation.StreamListener; import org.springframework.cloud.stream.messaging.Sink; import org.springframework.stereotype.Component; @Component public class ServiceBConsumer { @StreamListener(Sink.INPUT) public void handleMessage(String message) { System.out.println("Service B Instance received: " + message); } }
6. 测试生产者代码(可选)
import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.messaging.Source; import org.springframework.integration.support.MessageBuilder; import org.springframework.stereotype.Component; @Component @EnableBinding(Source.class) public class DataProducer { private final Source source; public DataProducer(Source source) { this.source = source; } public void sendMessage(String message) { source.output().send(MessageBuilder.withPayload(message).build()); } }
内容的提问来源于stack exchange,提问作者jarvis
相关产品推荐
相关产品推荐

