Spring Cloud Stream:当前连接异常时重发至另一AMQP连接方案咨询
问题2:异步监听发布结果(避免同步,覆写Bean实现)
完全可以通过覆写Spring Cloud Stream提供的Bean实现异步监听,不需要同步等待发送结果。核心是利用RabbitMQ的发布确认和消息退回机制,通过RabbitTemplateCustomizer配置回调:
先开启发布确认与退回:配置已经在问题1中包含,这里再强调下关键配置项:
spring: cloud: stream: rabbit: bindings: your-output-channel: producer: publisher-confirm-type: CORRELATED publisher-returns: true覆写RabbitTemplateCustomizer配置回调:
@Slf4j @Bean public RabbitTemplateCustomizer rabbitTemplateCustomizer() { return rabbitTemplate -> { // 发布确认回调:MQ收到消息后触发 rabbitTemplate.setConfirmCallback((correlationData, ack, cause) -> { if (correlationData != null) { if (ack) { log.info("消息[{}]已成功投递到MQ", correlationData.getId()); // 这里可以处理成功后的业务逻辑,比如更新消息状态 } else { log.error("消息[{}]投递失败,原因: {}", correlationData.getId(), cause); // 触发故障转移逻辑,比如切换到备用binder重发 } } }); // 消息退回回调:消息无法路由到队列时触发 rabbitTemplate.setReturnCallback((message, replyCode, replyText, exchange, routingKey) -> { log.error("消息[{}]路由失败,exchange: {}, routingKey: {}, 原因: {}", message.getMessageProperties().getCorrelationId(), exchange, routingKey, replyText); }); }; }发送消息时记得传入
CorrelationData关联消息与回调结果:@Autowired private MessageChannel yourOutputChannel; public void sendEvent(Object payload) { Message<?> message = MessageBuilder.withPayload(payload) .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON) .build(); // 传入唯一标识,用于回调时识别消息 yourOutputChannel.send(message, new CorrelationData(UUID.randomUUID().toString())); }这种方式完全异步,发送后无需等待结果,回调会在MQ返回确认/退回时自动触发,完美避免同步阻塞。
问题3:RabbitMQ HA是否有用?故障时的表现?
首先要明确:你当前是两台独立的RabbitMQ服务器,不是集群,所以RabbitMQ的HA队列(镜像队列)特性无法直接使用——HA队列是集群内的特性,需要节点加入同一集群才能配置镜像。
如果将两台服务器组成RabbitMQ集群并配置HA策略,那确实能大幅提升可靠性:
- HA队列会把队列镜像到集群多个节点,当承载队列的主节点故障时,集群会自动将某个从节点提升为主节点;只要客户端配置了集群地址列表(比如
spring.rabbitmq.addresses=node1:5672,node2:5672),就能自动重新连接到新主节点,发布和处理都不会崩溃。 - 你不需要保证消息顺序,可将HA策略的
sync-mode设为automatic,既保证数据同步,又不会有太大性能损耗;结合持久化队列,即使主节点故障,消息也不会丢失。
但如果保持当前的独立服务器架构,HA队列完全帮不上忙——两台服务器无集群关联,队列各自独立,一台故障后,另一台的队列结构虽相同但无数据同步,故障转移只能依赖Spring Cloud Stream的重试机制完成。
内容的提问来源于stack exchange,提问作者Roman Lebedev
相关产品推荐
相关产品推荐

