You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 01:20:26