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

如何避免Spring Cloud Dataflow Source重复返回RabbitMQ队列旧消息

问题

在Spring Cloud Dataflow Source应用中订阅RabbitMQ队列,消费消息以获取消息体和特定Header。当前应用启动后会持续执行并重复返回最后一条消息,需要调整配置或代码,让应用仅在新消息到达监听队列时执行。

现有Supplier代码

@Log4j2
@EnableConfigurationProperties({ RabbitMQProperties.class })
@Configuration
public class ReceiveMessageConfiguration {
    private final static String QUEUE_NAME = "UniversalId";
    String payload;
    Connection connection;
    Channel channel;


    @Bean
    public Supplier<String> receiverMessage(RabbitMQProperties rabbitMQProperties) throws Exception {
        return () -> {
            try {
                ConnectionFactory factory = new ConnectionFactory();
                factory.setHost("localhost");
                factory.setRequestedHeartbeat(0);

                connection = factory.newConnection();
                channel = connection.createChannel();

                channel.queueDeclare(QUEUE_NAME, true, false, false, null);

                DeliverCallback deliverCallback = (consumerTag, delivery) -> {
                    String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
//              Get headers from properties
                    AMQP.BasicProperties properties = delivery.getProperties();
                    Map<String, Object> headers = properties.getHeaders();

//              Extract and print payload and header
                    for (Map.Entry<String, Object> header : headers.entrySet()) {
                        if (header.getKey().toString().equals("UniversalId")) {
//                      log.info("ID nedeed: " + header.getValue());
                        }
                    }
                    payload = message;
                };
                channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> {
                });
                log.info(payload);
                if (payload != null) {
                    return payload;
                }

            } catch (Exception e) {
                e.printStackTrace();
            } finally {

                try {
                    channel.close();
                    connection.close();
                } catch (IOException | TimeoutException e) {
                    e.printStackTrace();
                }
            }
            return new String("No message");

        };
    }
}

配置文件

spring.cloud.stream.function.bindings.receiverMessage-in-0=input
spring.cloud.stream.function.bindings.receiverMessage-out-0=output
spring.cloud.stream.bindings.input.destination=receiverMessage-input
spring.cloud.stream.bindings.output.destination=receiverMessage-output

控制台日志示例

2022-12-21T16:25:19.725-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[           main][0;39m [36mm.c.n.r.RabbitMqCustomSourceApplication [0;39m [2m:[0;39m Started RabbitMqCustomSourceApplication in 7.888 seconds (process running for 9.482)
[2m2022-12-21T16:25:20.432-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m null
[2m2022-12-21T16:25:21.478-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m null
[2m2022-12-21T16:25:22.523-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m null
[2m2022-12-21T16:25:23.568-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m null
[2m2022-12-21T16:25:24.604-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m null
[2m2022-12-21T16:25:25.913-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m null
[2m2022-12-21T16:25:27.056-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m null
[2m2022-12-21T16:25:28.124-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m null
[2m2022-12-21T16:25:29.167-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m [
    {
        "MailRequest": {
            "action": "fourth",
            "mails": []
        }
    }
]
[2m2022-12-21T16:25:30.233-06:00[0;39m [32m INFO[0;39m [35m21600[0;39m [2m---[0;39m [2m[   scheduling-1][0;39m [36mm.c.n.r.s.ReceiveMessageConfiguration   [0;39m [2m:[0;39m [
    {
        "MailRequest": {
            "action": "fourth",
            "mails": []
        }
    }
]

已尝试从RabbitMQ API切换到AMQP API,但行为依旧。


解决方案

问题根源

  1. Supplier类型误用:Spring Cloud Stream中Supplier默认是轮询触发(默认间隔1秒),这就是控制台每秒输出一次结果、无新消息时重复返回旧payload的原因。
  2. 连接管理错误:每次Supplier执行都创建新RabbitMQ连接/通道,执行完就关闭,既浪费资源,也无法保持长连接监听新消息。
  3. 异步回调问题:basicConsume是异步执行的,Supplier可能在回调处理消息前就返回初始null或缓存的旧值。

修改方案

方案1:使用Spring Cloud Stream原生绑定(推荐)

无需手动管理RabbitMQ连接,直接利用框架的消息监听能力:

@Log4j2
@Configuration
public class ReceiveMessageConfiguration {

    @Bean
    public Consumer<Message<String>> receiverMessage() {
        return message -> {
            // 获取消息体
            String payload = message.getPayload();
            // 获取指定Header
            String universalId = message.getHeaders().get("UniversalId", String.class);
            
            log.info("收到消息体: {}", payload);
            log.info("获取到UniversalId: {}", universalId);
            
            // 若需输出到下游,可改为Function<Message<String>, String>并返回结果
        };
    }
}

对应配置调整:

# 指定绑定器为RabbitMQ
spring.cloud.stream.bindings.receiverMessage-in-0.binder=rabbit
# 直接绑定到目标队列
spring.cloud.stream.bindings.receiverMessage-in-0.destination=UniversalId
# 匹配队列持久化配置
spring.cloud.stream.rabbit.bindings.receiverMessage-in-0.consumer.durable-subscription=true
spring.cloud.stream.rabbit.bindings.receiverMessage-in-0.consumer.queue-name-group-only=true
# 保留原始消息格式
spring.cloud.stream.bindings.receiverMessage-in-0.content-type=application/json

方案2:手动实现长连接+阻塞消费(不推荐)

若必须手动管理连接,需保持长连接并通过阻塞队列等待新消息:

@Log4j2
@EnableConfigurationProperties({ RabbitMQProperties.class })
@Configuration
public class ReceiveMessageConfiguration {
    private final static String QUEUE_NAME = "UniversalId";
    private final BlockingQueue<String> messageQueue = new LinkedBlockingQueue<>();
    private Channel channel;

    // 初始化长连接与监听器
    @PostConstruct
    public void initRabbitMQ() throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        factory.setRequestedHeartbeat(60); // 设置心跳避免连接断开

        Connection connection = factory.newConnection();
        channel = connection.createChannel();
        channel.queueDeclare(QUEUE_NAME, true, false, false, null);

        DeliverCallback deliverCallback = (consumerTag, delivery) -> {
            String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
            String universalId = delivery.getProperties().getHeaders().get("UniversalId").toString();
            
            log.info("获取到UniversalId: {}", universalId);
            messageQueue.put(message); // 消息放入阻塞队列
        };
        channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> {});
        log.info("RabbitMQ监听器已启动,等待新消息...");
    }

    @Bean
    public Supplier<String> receiverMessage() {
        return () -> {
            try {
                return messageQueue.take(); // 阻塞等待新消息
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                return "消费中断";
            }
        };
    }

    // 销毁时关闭资源
    @PreDestroy
    public void cleanup() throws IOException, TimeoutException {
        if (channel != null && channel.isOpen()) {
            channel.close();
        }
    }
}

同时关闭Supplier轮询:

# 禁用自动轮询
spring.cloud.stream.poller.fixed-delay=-1

关键说明

  • 优先选择方案1,Spring Cloud Stream已封装RabbitMQ连接、监听等逻辑,减少重复代码与出错概率。
  • 方案2通过阻塞队列避免轮询重复返回旧消息,同时保持长连接实时监听新消息。

内容的提问来源于stack exchange,提问作者Emmanuel_Arreola

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 23:40:40