From 8a19df570914866ad62782817d9e9cfb97041f0f Mon Sep 17 00:00:00 2001 From: yinhuaiwei Date: Fri, 11 Sep 2026 09:36:03 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20=E4=BC=98=E5=8C=96=E7=9B=90?= =?UTF-8?q?=E6=BA=90ws=E6=8E=A8=E9=80=81=EF=BC=8C=E6=94=AF=E6=8C=81?= =?UTF-8?q?=E4=BB=BB=E6=84=8F=E6=B6=88=E6=81=AF=E6=A0=BC=E5=BC=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- inspect-main/inspect-main-api/pom.xml | 4 ++ .../inspect/api/config/RabbitMQConfig.java | 45 +++++++++++++++ .../api/constant/YanyuanConstants.java | 10 ++++ .../api/controller/AlarmController.java | 3 +- .../inspect/api/service/AlarmConsumer.java | 55 ++++++++++++++++--- .../inspect/api/service/AlarmPublisher.java | 7 +-- .../api/websocket/AlarmWsConnection.java | 6 -- 7 files changed, 112 insertions(+), 18 deletions(-) diff --git a/inspect-main/inspect-main-api/pom.xml b/inspect-main/inspect-main-api/pom.xml index 32fcdc5..7448d5e 100644 --- a/inspect-main/inspect-main-api/pom.xml +++ b/inspect-main/inspect-main-api/pom.xml @@ -45,6 +45,10 @@ com.inspect inspect-base-core + + com.inspect + inspect-base-redis + diff --git a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/config/RabbitMQConfig.java b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/config/RabbitMQConfig.java index 2700d94..da0f479 100644 --- a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/config/RabbitMQConfig.java +++ b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/config/RabbitMQConfig.java @@ -4,6 +4,7 @@ import com.inspect.api.constant.YanyuanConstants; import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.Queue; +import org.springframework.amqp.core.QueueBuilder; import org.springframework.amqp.core.TopicExchange; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.amqp.support.converter.MessageConverter; @@ -13,6 +14,9 @@ import org.springframework.context.annotation.Configuration; @Configuration public class RabbitMQConfig { + /** 告警重试延迟(毫秒),重试队列消息的 TTL */ + private static final int ALARM_RETRY_TTL_MS = 30_000; + @Bean public TopicExchange alarmExchange() { return new TopicExchange(YanyuanConstants.EXCHANGE, true, false); @@ -28,6 +32,47 @@ public class RabbitMQConfig { return BindingBuilder.bind(alarmQueue()).to(alarmExchange()).with(YanyuanConstants.ROUTING_KEY); } + // ============================ 告警重试/死信 ============================ + + @Bean + public TopicExchange alarmRetryExchange() { + return new TopicExchange(YanyuanConstants.ALARM_RETRY_EXCHANGE, true, false); + } + + /** 重试队列:消息 TTL 到期后经死信交换机自动投回主队列 */ + @Bean("alarmRetryQueue") + public Queue alarmRetryQueue() { + return QueueBuilder.durable(YanyuanConstants.ALARM_RETRY_QUEUE) + .ttl(ALARM_RETRY_TTL_MS) + .deadLetterExchange(YanyuanConstants.EXCHANGE) + .deadLetterRoutingKey(YanyuanConstants.ROUTING_KEY) + .build(); + } + + @Bean + public Binding alarmRetryBinding() { + return BindingBuilder.bind(alarmRetryQueue()) + .to(alarmRetryExchange()) + .with(YanyuanConstants.ALARM_RETRY_ROUTING_KEY); + } + + @Bean + public TopicExchange alarmDeadExchange() { + return new TopicExchange(YanyuanConstants.ALARM_DEAD_EXCHANGE, true, false); + } + + @Bean("alarmDeadQueue") + public Queue alarmDeadQueue() { + return QueueBuilder.durable(YanyuanConstants.ALARM_DEAD_QUEUE).build(); + } + + @Bean + public Binding alarmDeadBinding() { + return BindingBuilder.bind(alarmDeadQueue()) + .to(alarmDeadExchange()) + .with(YanyuanConstants.ALARM_DEAD_ROUTING_KEY); + } + // ============================ 巡检报告上传推送 MQ ============================ @Bean diff --git a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/constant/YanyuanConstants.java b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/constant/YanyuanConstants.java index 67bb27e..3a1d55c 100644 --- a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/constant/YanyuanConstants.java +++ b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/constant/YanyuanConstants.java @@ -8,6 +8,16 @@ public class YanyuanConstants { public static final String QUEUE = "alarm.queue"; public static final String ROUTING_KEY = "alarm.routing-key"; + /** 告警重试(WS发送失败后延迟重试) */ + public static final String ALARM_RETRY_EXCHANGE = "alarm.retry.exchange"; + public static final String ALARM_RETRY_QUEUE = "alarm.retry.queue"; + public static final String ALARM_RETRY_ROUTING_KEY = "alarm.retry.routing-key"; + + /** 告警死信(重试耗尽后兜底) */ + public static final String ALARM_DEAD_EXCHANGE = "alarm.dead.exchange"; + public static final String ALARM_DEAD_QUEUE = "alarm.dead.queue"; + public static final String ALARM_DEAD_ROUTING_KEY = "alarm.dead.routing-key"; + /** 巡检报告上传推送(先上报MQ,再由消费者HTTP上传外部平台) */ public static final String REPORT_EXCHANGE = "report.exchange"; public static final String REPORT_QUEUE = "report.upload.queue"; diff --git a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/controller/AlarmController.java b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/controller/AlarmController.java index d477b58..7111ea6 100644 --- a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/controller/AlarmController.java +++ b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/controller/AlarmController.java @@ -16,6 +16,7 @@ import org.springframework.web.bind.annotation.*; import java.util.ArrayList; import java.util.List; +import java.util.Map; @Slf4j @RestController @@ -52,7 +53,7 @@ public class AlarmController extends BaseController { @PostMapping("/push") @ResponseBody - public AjaxResult push(AlarmMessage msg) { + public AjaxResult push(@RequestBody Map msg) { if (msg != null) { alarmPublisher.publish(msg); return AjaxResult.success(msg); diff --git a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmConsumer.java b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmConsumer.java index f941e62..3839204 100644 --- a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmConsumer.java +++ b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmConsumer.java @@ -1,37 +1,78 @@ package com.inspect.api.service; +import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; -import com.inspect.api.domain.AlarmMessage; +import com.inspect.api.constant.YanyuanConstants; import com.inspect.api.websocket.AlarmWsConnection; +import com.inspect.base.core.utils.StringUtils; +import com.inspect.base.redis.service.RedisService; import lombok.extern.slf4j.Slf4j; +import org.springframework.amqp.core.Message; import org.springframework.amqp.rabbit.annotation.RabbitListener; +import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Component; +import java.util.Map; + @Slf4j @Component public class AlarmConsumer { + + /** 最多重试次数(总尝试次数 = MAX_RETRY + 1) */ + private static final int MAX_RETRY = 2; + private static final String RETRY_COUNT_HEADER = "alarm-retry-count"; + private final AlarmWsConnection wsConnection; private final ObjectMapper objectMapper; + private final RedisService redisService; + private final RabbitTemplate rabbitTemplate; - public AlarmConsumer(AlarmWsConnection wsConnection, ObjectMapper objectMapper) { + public AlarmConsumer(AlarmWsConnection wsConnection, ObjectMapper objectMapper, + RedisService redisService, RabbitTemplate rabbitTemplate) { this.wsConnection = wsConnection; this.objectMapper = objectMapper; + this.redisService = redisService; + this.rabbitTemplate = rabbitTemplate; } @RabbitListener(queues = "#{alarmQueue.name}") - public void onAlarm(AlarmMessage msg) { + public void onAlarm(Message message) { try { + Map msg = objectMapper.readValue(message.getBody(), + new TypeReference>() {}); + String staticCode = (String) redisService.redisTemplate.opsForValue().get("STATION_CODE"); + if (StringUtils.isNotEmpty(staticCode)) { + msg.put("stationCode", staticCode); + } String json = objectMapper.writeValueAsString(msg); boolean sent = wsConnection.send(json); if (sent) { log.info("Alarm sent to WS: msg={}", msg); } else { - log.warn("Alarm not sent, WS not connected: msg={}", msg); - throw new RuntimeException("WS not connected, message will be requeued"); + retryOrDeadLetter(message); } } catch (Exception e) { - log.error("Alarm consumer error: {}", e.getMessage()); - throw new RuntimeException("Alarm processing failed", e); + log.error("Alarm consumer error: {}", e.getMessage(), e); + retryOrDeadLetter(message); } } + + private void retryOrDeadLetter(Message message) { + int retryCount = getRetryCount(message); + if (retryCount < MAX_RETRY) { + message.getMessageProperties().setHeader(RETRY_COUNT_HEADER, retryCount + 1); + rabbitTemplate.send(YanyuanConstants.ALARM_RETRY_EXCHANGE, + YanyuanConstants.ALARM_RETRY_ROUTING_KEY, message); + log.warn("Alarm WS send failed, retry {}/{} scheduled", retryCount + 1, MAX_RETRY); + } else { + rabbitTemplate.send(YanyuanConstants.ALARM_DEAD_EXCHANGE, + YanyuanConstants.ALARM_DEAD_ROUTING_KEY, message); + log.error("Alarm WS send failed after {} retries, moved to dead letter queue", MAX_RETRY); + } + } + + private int getRetryCount(Message message) { + Object count = message.getMessageProperties().getHeader(RETRY_COUNT_HEADER); + return count instanceof Number ? ((Number) count).intValue() : 0; + } } diff --git a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmPublisher.java b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmPublisher.java index ac38275..dfeadf6 100644 --- a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmPublisher.java +++ b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmPublisher.java @@ -1,7 +1,6 @@ package com.inspect.api.service; import com.inspect.api.constant.YanyuanConstants; -import com.inspect.api.domain.AlarmMessage; import lombok.extern.slf4j.Slf4j; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Service; @@ -15,12 +14,12 @@ public class AlarmPublisher { this.rabbitTemplate = rabbitTemplate; } - public void publish(AlarmMessage msg) { + public void publish(Object msg) { try { - rabbitTemplate.convertAndSend(YanyuanConstants.EXCHANGE, YanyuanConstants.ROUTING_KEY, msg); log.info("Alert published: msg={}", msg); + rabbitTemplate.convertAndSend(YanyuanConstants.EXCHANGE, YanyuanConstants.ROUTING_KEY, msg); } catch (Exception e) { - log.error("Failed to publish alert: {}", e.getMessage()); + log.error("Failed to publish alert: {}", e.getMessage(), e); } } } diff --git a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/websocket/AlarmWsConnection.java b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/websocket/AlarmWsConnection.java index e47e405..15ea0f1 100644 --- a/inspect-main/inspect-main-api/src/main/java/com/inspect/api/websocket/AlarmWsConnection.java +++ b/inspect-main/inspect-main-api/src/main/java/com/inspect/api/websocket/AlarmWsConnection.java @@ -13,7 +13,6 @@ import org.springframework.web.socket.WebSocketSession; import org.springframework.web.socket.client.standard.StandardWebSocketClient; import org.springframework.web.socket.handler.TextWebSocketHandler; -import javax.annotation.PostConstruct; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; @@ -34,11 +33,6 @@ public class AlarmWsConnection { this.wsClient = new StandardWebSocketClient(); } - @PostConstruct - public void init() { - connect(); - } - public void connect() { if (connecting) { return;