如何避免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,但行为依旧。
解决方案
问题根源
- Supplier类型误用:Spring Cloud Stream中
Supplier默认是轮询触发(默认间隔1秒),这就是控制台每秒输出一次结果、无新消息时重复返回旧payload的原因。 - 连接管理错误:每次Supplier执行都创建新RabbitMQ连接/通道,执行完就关闭,既浪费资源,也无法保持长连接监听新消息。
- 异步回调问题:
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
相关产品推荐
相关产品推荐

