Browse Source

refactor: 优化盐源ws推送,支持任意消息格式

yanyuan
yinhuaiwei 2 weeks ago
parent
commit
8a19df5709
7 changed files with 112 additions and 18 deletions
  1. +4
    -0
      inspect-main/inspect-main-api/pom.xml
  2. +45
    -0
      inspect-main/inspect-main-api/src/main/java/com/inspect/api/config/RabbitMQConfig.java
  3. +10
    -0
      inspect-main/inspect-main-api/src/main/java/com/inspect/api/constant/YanyuanConstants.java
  4. +2
    -1
      inspect-main/inspect-main-api/src/main/java/com/inspect/api/controller/AlarmController.java
  5. +48
    -7
      inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmConsumer.java
  6. +3
    -4
      inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmPublisher.java
  7. +0
    -6
      inspect-main/inspect-main-api/src/main/java/com/inspect/api/websocket/AlarmWsConnection.java

+ 4
- 0
inspect-main/inspect-main-api/pom.xml View File

@ -45,6 +45,10 @@
<groupId>com.inspect</groupId> <groupId>com.inspect</groupId>
<artifactId>inspect-base-core</artifactId> <artifactId>inspect-base-core</artifactId>
</dependency> </dependency>
<dependency>
<groupId>com.inspect</groupId>
<artifactId>inspect-base-redis</artifactId>
</dependency>
</dependencies> </dependencies>
<dependencyManagement> <dependencyManagement>
<dependencies> <dependencies>


+ 45
- 0
inspect-main/inspect-main-api/src/main/java/com/inspect/api/config/RabbitMQConfig.java View File

@ -4,6 +4,7 @@ import com.inspect.api.constant.YanyuanConstants;
import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.Queue; import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.QueueBuilder;
import org.springframework.amqp.core.TopicExchange; import org.springframework.amqp.core.TopicExchange;
import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.amqp.support.converter.MessageConverter;
@ -13,6 +14,9 @@ import org.springframework.context.annotation.Configuration;
@Configuration @Configuration
public class RabbitMQConfig { public class RabbitMQConfig {
/** 告警重试延迟(毫秒),重试队列消息的 TTL */
private static final int ALARM_RETRY_TTL_MS = 30_000;
@Bean @Bean
public TopicExchange alarmExchange() { public TopicExchange alarmExchange() {
return new TopicExchange(YanyuanConstants.EXCHANGE, true, false); return new TopicExchange(YanyuanConstants.EXCHANGE, true, false);
@ -28,6 +32,47 @@ public class RabbitMQConfig {
return BindingBuilder.bind(alarmQueue()).to(alarmExchange()).with(YanyuanConstants.ROUTING_KEY); 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 ============================ // ============================ 巡检报告上传推送 MQ ============================
@Bean @Bean


+ 10
- 0
inspect-main/inspect-main-api/src/main/java/com/inspect/api/constant/YanyuanConstants.java View File

@ -8,6 +8,16 @@ public class YanyuanConstants {
public static final String QUEUE = "alarm.queue"; public static final String QUEUE = "alarm.queue";
public static final String ROUTING_KEY = "alarm.routing-key"; 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上传外部平台) */ /** 巡检报告上传推送(先上报MQ,再由消费者HTTP上传外部平台) */
public static final String REPORT_EXCHANGE = "report.exchange"; public static final String REPORT_EXCHANGE = "report.exchange";
public static final String REPORT_QUEUE = "report.upload.queue"; public static final String REPORT_QUEUE = "report.upload.queue";


+ 2
- 1
inspect-main/inspect-main-api/src/main/java/com/inspect/api/controller/AlarmController.java View File

@ -16,6 +16,7 @@ import org.springframework.web.bind.annotation.*;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map;
@Slf4j @Slf4j
@RestController @RestController
@ -52,7 +53,7 @@ public class AlarmController extends BaseController {
@PostMapping("/push") @PostMapping("/push")
@ResponseBody @ResponseBody
public AjaxResult push(AlarmMessage msg) {
public AjaxResult push(@RequestBody Map<String, Object> msg) {
if (msg != null) { if (msg != null) {
alarmPublisher.publish(msg); alarmPublisher.publish(msg);
return AjaxResult.success(msg); return AjaxResult.success(msg);


+ 48
- 7
inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmConsumer.java View File

@ -1,37 +1,78 @@
package com.inspect.api.service; package com.inspect.api.service;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper; 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.api.websocket.AlarmWsConnection;
import com.inspect.base.core.utils.StringUtils;
import com.inspect.base.redis.service.RedisService;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import java.util.Map;
@Slf4j @Slf4j
@Component @Component
public class AlarmConsumer { 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 AlarmWsConnection wsConnection;
private final ObjectMapper objectMapper; 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.wsConnection = wsConnection;
this.objectMapper = objectMapper; this.objectMapper = objectMapper;
this.redisService = redisService;
this.rabbitTemplate = rabbitTemplate;
} }
@RabbitListener(queues = "#{alarmQueue.name}") @RabbitListener(queues = "#{alarmQueue.name}")
public void onAlarm(AlarmMessage msg) {
public void onAlarm(Message message) {
try { try {
Map<String, Object> msg = objectMapper.readValue(message.getBody(),
new TypeReference<Map<String, Object>>() {});
String staticCode = (String) redisService.redisTemplate.opsForValue().get("STATION_CODE");
if (StringUtils.isNotEmpty(staticCode)) {
msg.put("stationCode", staticCode);
}
String json = objectMapper.writeValueAsString(msg); String json = objectMapper.writeValueAsString(msg);
boolean sent = wsConnection.send(json); boolean sent = wsConnection.send(json);
if (sent) { if (sent) {
log.info("Alarm sent to WS: msg={}", msg); log.info("Alarm sent to WS: msg={}", msg);
} else { } 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) { } 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;
}
} }

+ 3
- 4
inspect-main/inspect-main-api/src/main/java/com/inspect/api/service/AlarmPublisher.java View File

@ -1,7 +1,6 @@
package com.inspect.api.service; package com.inspect.api.service;
import com.inspect.api.constant.YanyuanConstants; import com.inspect.api.constant.YanyuanConstants;
import com.inspect.api.domain.AlarmMessage;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
@ -15,12 +14,12 @@ public class AlarmPublisher {
this.rabbitTemplate = rabbitTemplate; this.rabbitTemplate = rabbitTemplate;
} }
public void publish(AlarmMessage msg) {
public void publish(Object msg) {
try { try {
rabbitTemplate.convertAndSend(YanyuanConstants.EXCHANGE, YanyuanConstants.ROUTING_KEY, msg);
log.info("Alert published: msg={}", msg); log.info("Alert published: msg={}", msg);
rabbitTemplate.convertAndSend(YanyuanConstants.EXCHANGE, YanyuanConstants.ROUTING_KEY, msg);
} catch (Exception e) { } catch (Exception e) {
log.error("Failed to publish alert: {}", e.getMessage());
log.error("Failed to publish alert: {}", e.getMessage(), e);
} }
} }
} }

+ 0
- 6
inspect-main/inspect-main-api/src/main/java/com/inspect/api/websocket/AlarmWsConnection.java View File

@ -13,7 +13,6 @@ import org.springframework.web.socket.WebSocketSession;
import org.springframework.web.socket.client.standard.StandardWebSocketClient; import org.springframework.web.socket.client.standard.StandardWebSocketClient;
import org.springframework.web.socket.handler.TextWebSocketHandler; import org.springframework.web.socket.handler.TextWebSocketHandler;
import javax.annotation.PostConstruct;
import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.atomic.AtomicReference;
@ -34,11 +33,6 @@ public class AlarmWsConnection {
this.wsClient = new StandardWebSocketClient(); this.wsClient = new StandardWebSocketClient();
} }
@PostConstruct
public void init() {
connect();
}
public void connect() { public void connect() {
if (connecting) { if (connecting) {
return; return;


Loading…
Cancel
Save