Revert "冀北版本过滤无效数据"

This reverts commit c1a2f3e8fb.
This commit is contained in:
hzj
2026-07-22 16:04:57 +08:00
parent c1a2f3e8fb
commit 14b213fb99
23 changed files with 1004 additions and 349 deletions

View File

@@ -12,7 +12,6 @@ import com.njcn.algorithm.pojo.api.fallback.PqReasonableRangeFeignClientFallback
import com.njcn.algorithm.pojo.dto.PqReasonableRangeDto;
import com.njcn.algorithm.pojo.param.DataCleanParam;
import com.njcn.common.pojo.constant.ServerInfo;
import com.njcn.common.pojo.response.HttpResult;
import io.swagger.annotations.ApiOperation;
import org.springframework.cloud.openfeign.FeignClient;

View File

@@ -20,7 +20,7 @@ import java.time.LocalDateTime;
@Data
public class MessageHarmonicDataSet implements Serializable {
private Integer FLAG;
@JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss.SSS")
@JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss")
@JsonDeserialize(using = LocalDateTimeDeserializer.class)
@JsonSerialize(using = LocalDateTimeSerializer.class)
private LocalDateTime TIME;

View File

@@ -25,7 +25,6 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.time.temporal.ChronoField;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
@@ -94,7 +93,7 @@ public class MessageAnalysisServiceImpl implements MessageAnalysisService {
MessageHarmonicDataSet messageHarmonicDataSet = JSONObject.parseObject(value, MessageHarmonicDataSet.class);
LocalDateTime localDateTime = messageHarmonicDataSet.getTIME();
//排除上电下电等情况前置上送上不是整分的数据
if((localDateTime.getSecond() != 0)||localDateTime.get(ChronoField.MILLI_OF_SECOND)!=0){
if(!(localDateTime.getSecond() == 0)){
return;
}
Integer flag = messageHarmonicDataSet.getFLAG();

View File

@@ -28,7 +28,6 @@ import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.*;
import java.time.LocalDateTime;
import java.time.temporal.ChronoField;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
@@ -87,10 +86,7 @@ public class LnDataDealController extends BaseController {
if(Objects.equals(DataTypeEnum.HARMONIC.getCode(),dataType)) {
MessageHarmonicDataSet messageHarmonicDataSet = JSONObject.parseObject(value, MessageHarmonicDataSet.class);
LocalDateTime localDateTime = messageHarmonicDataSet.getTIME();
//排除上电下电等情况前置上送上不是整分的数据
if((localDateTime.getSecond() != 0)||localDateTime.get(ChronoField.MILLI_OF_SECOND)!=0){
System.out.println(111111111);
}
Integer flag = messageHarmonicDataSet.getFLAG();
MessageP pq = messageHarmonicDataSet.getPQ();
MessageV v = messageHarmonicDataSet.getV();

View File

@@ -105,12 +105,12 @@
<artifactId>mq-spring-boot-starter</artifactId>
<version>1.0.0-SNAPSHOT</version>
</dependency>
<!-- <dependency>-->
<!-- <groupId>org.springframework.boot</groupId>-->
<!-- <artifactId>spring-boot-starter-test</artifactId>-->
<!-- <version>2.7.12</version>-->
<!-- <scope>test</scope>-->
<!-- </dependency>-->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<version>2.7.12</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>

View File

@@ -6,7 +6,6 @@ import com.njcn.message.messagedto.DevComFlagDTO;
import com.njcn.message.constant.RedisKeyPrefix;
import com.njcn.middle.rocket.constant.EnhanceMessageConstant;
import com.njcn.middle.rocket.handler.EnhanceConsumerMessageHandler;
import com.njcn.mq.annotation.MqListener;
import com.njcn.redis.pojo.enums.AppRedisKey;
import com.njcn.redis.pojo.enums.RedisKeyEnum;
import com.njcn.redis.utils.RedisUtil;
@@ -31,18 +30,59 @@ import java.util.Objects;
* @version V1.0.0
*/
@Component
@RocketMQMessageListener(
topic = "Device_Run_Flag_Topic",
consumerGroup = "Device_Run_Flag_Consumer",
selectorExpression = "*",
consumeThreadNumber = 10,
enableMsgTrace = true
)
@Slf4j
public class DeviceRunFlagDataConsumer {
public class DeviceRunFlagDataConsumer extends EnhanceConsumerMessageHandler<DevComFlagDTO> implements RocketMQListener<String> {
@Resource
private RedisUtil redisUtil;
@Resource
private RocketMqLogFeignClient rocketMqLogFeignClient;
@Autowired
private MessAnalysisFeignClient messAnalysisFeignClient;
@Override
public void onMessage(String message) {
DevComFlagDTO devComFlagDTO = JSONObject.parseObject(message,DevComFlagDTO.class);
super.dispatchMessage(devComFlagDTO);
@MqListener(topic = "Device_Run_Flag_Topic", group = "Device_Run_Flag_Consumer")
protected void onMessage(DevComFlagDTO message) {
}
/***
* 通过redis分布式锁判断当前消息所处状态
* 1、null 查不到该key的数据属于第一次消费放行
* 2、fail 上次消息消费时发生异常,放行
* 3、being processed 正在处理,打回去
* 4、success 最近72小时消费成功避免重复消费打回去
*/
@Override
public boolean filter(DevComFlagDTO message) {
// String keyStatus = redisUtil.getStringByKey(AppRedisKey.RMQ_CONSUME_KEY.concat(message.getKey()));
// if (Objects.isNull(keyStatus) || keyStatus.equalsIgnoreCase(MessageStatus.FAIL)) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.DEVICE_RUN_FLAG.concat(message.getKey()), MessageStatus.BEING_PROCESSED, 30L);
// return false;
// }
return false;
}
/**
* 消费成功缓存到redis72小时避免重复消费
*/
@Override
protected void consumeSuccess(DevComFlagDTO message) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.DEVICE_RUN_FLAG.concat(message.getKey()), MessageStatus.SUCCESS, 5*60L);
}
@Override
protected void handleMessage(DevComFlagDTO message) {
//获取之前设备状态
//删除设备时前置会在连接一次通道但是DevId为空所以添加
if(StringUtils.isNoneBlank(message.getId())){
@@ -62,6 +102,53 @@ public class DeviceRunFlagDataConsumer {
/**
* 发生异常时,进行错误信息入库保存
* 默认没有实现类子类可以实现该方法调用feign接口进行入库保存
*/
@Override
protected void saveExceptionMsgLog(DevComFlagDTO message, String identity, Exception exception) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.DEVICE_RUN_FLAG.concat(message.getKey()), MessageStatus.FAIL, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
RocketmqMsgErrorLog rocketmqMsgErrorLog = new RocketmqMsgErrorLog();
rocketmqMsgErrorLog.setMsgKey(message.getKey());
rocketmqMsgErrorLog.setResource(message.getSource());
if (identity.equalsIgnoreCase(EnhanceMessageConstant.IDENTITY_SINGLE)) {
//数据库字段配置长度200避免插入失败大致分析异常原因
String exceptionMsg = exception.getMessage();
if(exceptionMsg.length() > 200){
exceptionMsg = exceptionMsg.substring(0,180);
}
rocketmqMsgErrorLog.setRecord(exceptionMsg);
//如果是当前消息重试的则略过
if(!message.getSource().startsWith(EnhanceMessageConstant.RETRY_PREFIX)){
//单次消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
} else {
rocketmqMsgErrorLog.setRecord("重试消费" + super.getMaxRetryTimes() + "次,依旧消费失败。");
//重试N次后依然消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
}
/***
* 处理失败后,是否重试
* 一般开启
*/
@Override
protected boolean isRetry() {
return true;
}
/***
* 消费失败是否抛出异常,抛出异常后就不再消费了
*/
@Override
protected boolean throwException() {
return false;
}

View File

@@ -8,7 +8,6 @@ import com.njcn.message.messagedto.MessageDataDTO;
import com.njcn.message.websocket.WebSocketServer;
import com.njcn.middle.rocket.constant.EnhanceMessageConstant;
import com.njcn.middle.rocket.handler.EnhanceConsumerMessageHandler;
import com.njcn.mq.annotation.MqListener;
import com.njcn.redis.pojo.enums.RedisKeyEnum;
import com.njcn.redis.utils.RedisUtil;
import com.njcn.system.api.RocketMqLogFeignClient;
@@ -32,17 +31,59 @@ import java.util.concurrent.TimeUnit;
* @version V1.0.0
*/
@Component
@RocketMQMessageListener(
topic = "File_Reply_Topic",
consumerGroup = "file_sys_consumer",
consumeThreadNumber = 10,
enableMsgTrace = true
)
@Slf4j
public class FileSysDataConsumer {
public class FileSysDataConsumer extends EnhanceConsumerMessageHandler<FileSysDTO> implements RocketMQListener<String> {
@Resource
private RedisUtil redisUtil;
@Resource
private StringRedisTemplate stringRedisTemplate;
@Resource
private RocketMqLogFeignClient rocketMqLogFeignClient;
@Override
public void onMessage(String message) {
@MqListener(topic = "File_Reply_Topic", group = "file_sys_consumer")
FileSysDTO messageDataDTO = JSONObject.parseObject(message, FileSysDTO.class);
super.dispatchMessage(messageDataDTO);
}
/***
* 通过redis分布式锁判断当前消息所处状态
* 1、null 查不到该key的数据属于第一次消费放行
* 2、fail 上次消息消费时发生异常,放行
* 3、being processed 正在处理,打回去
* 4、success 最近72小时消费成功避免重复消费打回去
*/
@Override
public boolean filter(FileSysDTO message) {
String keyStatus = redisUtil.getStringByKey(RedisKeyPrefix.REAL_TIME_DATA.concat(message.getKey()));
if (Objects.isNull(keyStatus) || keyStatus.equalsIgnoreCase(MessageStatus.FAIL)) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.REAL_TIME_DATA.concat(message.getKey()), MessageStatus.BEING_PROCESSED, 30L);
return false;
}
return true;
}
/**
* 消费成功缓存到redis5分钟避免重复消费
*/
@Override
protected void consumeSuccess(FileSysDTO message) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.REAL_TIME_DATA.concat(message.getKey()), MessageStatus.SUCCESS, 5*60L);
}
@Override
protected void handleMessage(FileSysDTO message) {
String msgId = message.getGuid();
@@ -56,8 +97,54 @@ public class FileSysDataConsumer {
/**
* 发生异常时,进行错误信息入库保存
* 默认没有实现类子类可以实现该方法调用feign接口进行入库保存
*/
@Override
protected void saveExceptionMsgLog(FileSysDTO message, String identity, Exception exception) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.REAL_TIME_DATA.concat(message.getKey()), MessageStatus.FAIL, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
RocketmqMsgErrorLog rocketmqMsgErrorLog = new RocketmqMsgErrorLog();
rocketmqMsgErrorLog.setMsgKey(message.getKey());
rocketmqMsgErrorLog.setResource(message.getSource());
if (identity.equalsIgnoreCase(EnhanceMessageConstant.IDENTITY_SINGLE)) {
//数据库字段配置长度200避免插入失败大致分析异常原因
String exceptionMsg = exception.getMessage();
if(exceptionMsg.length() > 200){
exceptionMsg = exceptionMsg.substring(0,180);
}
rocketmqMsgErrorLog.setRecord(exceptionMsg);
//如果是当前消息重试的则略过
if(!message.getSource().startsWith(EnhanceMessageConstant.RETRY_PREFIX)){
//单次消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
} else {
rocketmqMsgErrorLog.setRecord("重试消费" + super.getMaxRetryTimes() + "次,依旧消费失败。");
//重试N次后依然消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
}
/***
* 处理失败后,是否重试
* 一般开启
*/
@Override
protected boolean isRetry() {
return true;
}
/***
* 消费失败是否抛出异常,抛出异常后就不再消费了
*/
@Override
protected boolean throwException() {
return false;
}

View File

@@ -1,18 +1,34 @@
package com.njcn.message.consumer;
import com.alibaba.fastjson.JSONObject;
import com.njcn.common.pojo.enums.response.CommonResponseEnum;
import com.njcn.common.pojo.response.HttpResult;
import com.njcn.message.constant.MessageStatus;
import com.njcn.message.messagedto.MessageDataDTO;
import com.njcn.mq.annotation.MqListener;
import com.njcn.message.constant.RedisKeyPrefix;
import com.njcn.middle.rocket.constant.EnhanceMessageConstant;
import com.njcn.middle.rocket.handler.EnhanceConsumerMessageHandler;
import com.njcn.redis.pojo.enums.RedisKeyEnum;
import com.njcn.redis.utils.RedisUtil;
import com.njcn.stat.api.MessAnalysisFeignClient;
import com.njcn.system.api.RocketMqLogFeignClient;
import com.njcn.system.pojo.po.RocketmqMsgErrorLog;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import javax.annotation.PostConstruct;
import javax.annotation.Resource;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
/**
@@ -23,20 +39,151 @@ import java.util.List;
* @createTime 2023/8/11 15:32
*/
@Component
@org.springframework.boot.autoconfigure.condition.ConditionalOnProperty(
name = "mq.type", havingValue = "rocketmq", matchIfMissing = true)
@RocketMQMessageListener(
topic = "LN_Topic",
consumerGroup = "ln_consumer",
selectorExpression = "*",
consumeThreadNumber = 10,
enableMsgTrace = true
)
@Slf4j
public class FrontDataConsumer {
public class FrontDataConsumer extends EnhanceConsumerMessageHandler<MessageDataDTO> implements RocketMQListener<String> {
@Autowired
private MessAnalysisFeignClient messAnalysisFeignClient;
@MqListener(topic = "MQ_INST_1697676490574298_Bm3DLamM%LN_Topic", group = "ln_consumer",batchSize = 50,batchTimeoutMs=10000)
public void onMessage(List<MessageDataDTO> msg) {
log.info("开始消费");
messAnalysisFeignClient.analysis(msg);
@Value("${rocketmq.consumer_size}")
private Integer consumerSize;
@Resource
private RedisUtil redisUtil;
@Resource
private RocketMqLogFeignClient rocketMqLogFeignClient;
private List<MessageDataDTO> messageList = new ArrayList<>();
@PostConstruct
public void validateConfig() {
if (consumerSize == null) {
throw new IllegalStateException("rocketmq.consumer_size 未配置!");
}
this.messageList = new ArrayList<>(consumerSize);
}
@Override
public void onMessage(String baseMessage) {
MessageDataDTO messageDataDTO = JSONObject.parseObject(baseMessage,MessageDataDTO.class);
super.dispatchMessage(messageDataDTO);
}
/***
* 通过redis分布式锁判断当前消息所处状态
* 1、null 查不到该key的数据属于第一次消费放行
* 2、fail 上次消息消费时发生异常,放行
* 3、being processed 正在处理,打回去
* 4、success 最近72小时消费成功避免重复消费打回去
*/
@Override
public boolean filter(MessageDataDTO message) {
String keyStatus = redisUtil.getStringByKey(message.getKey());
if (Objects.isNull(keyStatus) || keyStatus.equalsIgnoreCase(MessageStatus.FAIL)) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.HARMMONIC_TOPIC.concat(message.getKey()), MessageStatus.BEING_PROCESSED, 60L);
return false;
}
return true;
}
/**
* 消费成功缓存到redis5分钟避免重复消费
*/
@Override
protected void consumeSuccess(MessageDataDTO message) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.HARMMONIC_TOPIC.concat(message.getKey()), MessageStatus.SUCCESS, 5*60L);
}
@Override
protected void handleMessage(MessageDataDTO message) {
synchronized (messageList) {
messageList.add(message);
if (messageList.size() >= consumerSize) {
saveToDatabase();
}
}
}
/**
* 发生异常时,进行错误信息入库保存
* 默认没有实现类子类可以实现该方法调用feign接口进行入库保存
*/
@Override
protected void saveExceptionMsgLog(MessageDataDTO message, String identity, Exception exception) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.HARMMONIC_TOPIC.concat(message.getKey()), MessageStatus.FAIL, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
RocketmqMsgErrorLog rocketmqMsgErrorLog = new RocketmqMsgErrorLog();
rocketmqMsgErrorLog.setMsgKey(message.getKey());
rocketmqMsgErrorLog.setResource(message.getSource());
if (identity.equalsIgnoreCase(EnhanceMessageConstant.IDENTITY_SINGLE)) {
//数据库字段配置长度200避免插入失败大致分析异常原因
String exceptionMsg = exception.getMessage();
if(exceptionMsg.length() > 200){
exceptionMsg = exceptionMsg.substring(0,180);
}
rocketmqMsgErrorLog.setRecord(exceptionMsg);
//如果是当前消息重试的则略过
if(!message.getSource().startsWith(EnhanceMessageConstant.RETRY_PREFIX)){
//单次消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
} else {
rocketmqMsgErrorLog.setRecord("重试消费" + super.getMaxRetryTimes() + "次,依旧消费失败。");
//重试N次后依然消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
}
/***
* 处理失败后,是否重试
* 一般开启
*/
@Override
protected boolean isRetry() {
return false;
}
/***
* 消费失败是否抛出异常,抛出异常后就不再消费了
*/
@Override
protected boolean throwException() {
return false;
}
//50个消息做一组插入数据库
public void saveToDatabase(){
try {
long start = System.currentTimeMillis();
messAnalysisFeignClient.analysis(messageList);
long end = System.currentTimeMillis();
log.info("处理"+consumerSize+"条消息所需时间------------"+(end-start));
}catch (Exception e){{
log.info(e.toString());
}
}finally{
messageList.clear();
}
}
}

View File

@@ -7,7 +7,6 @@ import com.njcn.message.messagedto.FrontHeartBeatDTO;
import com.njcn.message.constant.RedisKeyPrefix;
import com.njcn.middle.rocket.constant.EnhanceMessageConstant;
import com.njcn.middle.rocket.handler.EnhanceConsumerMessageHandler;
import com.njcn.mq.annotation.MqListener;
import com.njcn.redis.pojo.enums.AppRedisKey;
import com.njcn.redis.pojo.enums.RedisKeyEnum;
import com.njcn.redis.utils.RedisUtil;
@@ -31,17 +30,58 @@ import java.util.Objects;
* @version V1.0.0
*/
@Component
@RocketMQMessageListener(
topic = "Heart_Beat_Topic",
consumerGroup = "Heartb_Beat_Topic_Consumer",
selectorExpression = "*",
consumeThreadNumber = 10,
enableMsgTrace = true
)
@Slf4j
public class FrontHeartBeatConsumer {
public class FrontHeartBeatConsumer extends EnhanceConsumerMessageHandler<FrontHeartBeatDTO> implements RocketMQListener<String> {
@Resource
private RedisUtil redisUtil;
@Resource
private RocketMqLogFeignClient rocketMqLogFeignClient;
@Autowired
private MessAnalysisFeignClient messAnalysisFeignClient;
@MqListener(topic = "Heart_Beat_Topic", group = "Heartb_Beat_Topic_Consumer")
@Override
public void onMessage(String message) {
FrontHeartBeatDTO frontHeartBeatDTO = JSONObject.parseObject(message, FrontHeartBeatDTO.class);
super.dispatchMessage(frontHeartBeatDTO);
}
//本消息不需要控制重复消费
/***
* 通过redis分布式锁判断当前消息所处状态
* 1、null 查不到该key的数据属于第一次消费放行
* 2、fail 上次消息消费时发生异常,放行
* 3、being processed 正在处理,打回去
* 4、success 最近72小时消费成功避免重复消费打回去
*/
@Override
public boolean filter(FrontHeartBeatDTO message) {
// String keyStatus = redisUtil.getStringByKey(AppRedisKey.RMQ_CONSUME_KEY.concat(message.getKey()));
// if (Objects.isNull(keyStatus) || keyStatus.equalsIgnoreCase(MessageStatus.FAIL)) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.HEART_BEAT.concat(message.getKey()), MessageStatus.BEING_PROCESSED, 30L);
// return false;
// }
return false;
}
/**
* 消费成功缓存到redis72小时避免重复消费
*/
@Override
protected void consumeSuccess(FrontHeartBeatDTO message) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.HEART_BEAT.concat(message.getKey()), MessageStatus.SUCCESS, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
}
@Override
protected void handleMessage(FrontHeartBeatDTO message) {
//将心跳状态存到redis失效时间为30s如果持续有心跳则进程在线反之redis找不到则不在线
if(Objects.equals(message.getFronttype(), FrontTypeEnum.STAT.getCode())){
@@ -53,8 +93,54 @@ public class FrontHeartBeatConsumer {
/**
* 发生异常时,进行错误信息入库保存
* 默认没有实现类子类可以实现该方法调用feign接口进行入库保存
*/
@Override
protected void saveExceptionMsgLog(FrontHeartBeatDTO message, String identity, Exception exception) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.HEART_BEAT.concat(message.getKey()), MessageStatus.FAIL, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
RocketmqMsgErrorLog rocketmqMsgErrorLog = new RocketmqMsgErrorLog();
rocketmqMsgErrorLog.setMsgKey(message.getKey());
rocketmqMsgErrorLog.setResource(message.getSource());
if (identity.equalsIgnoreCase(EnhanceMessageConstant.IDENTITY_SINGLE)) {
//数据库字段配置长度200避免插入失败大致分析异常原因
String exceptionMsg = exception.getMessage();
if(exceptionMsg.length() > 200){
exceptionMsg = exceptionMsg.substring(0,180);
}
rocketmqMsgErrorLog.setRecord(exceptionMsg);
//如果是当前消息重试的则略过
if(!message.getSource().startsWith(EnhanceMessageConstant.RETRY_PREFIX)){
//单次消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
} else {
rocketmqMsgErrorLog.setRecord("重试消费" + super.getMaxRetryTimes() + "次,依旧消费失败。");
//重试N次后依然消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
}
/***
* 处理失败后,是否重试
* 一般开启
*/
@Override
protected boolean isRetry() {
return true;
}
/***
* 消费失败是否抛出异常,抛出异常后就不再消费了
*/
@Override
protected boolean throwException() {
return false;
}

View File

@@ -8,7 +8,6 @@ import com.njcn.message.constant.RedisKeyPrefix;
import com.njcn.message.websocket.WebSocketServer;
import com.njcn.middle.rocket.constant.EnhanceMessageConstant;
import com.njcn.middle.rocket.handler.EnhanceConsumerMessageHandler;
import com.njcn.mq.annotation.MqListener;
import com.njcn.redis.pojo.enums.RedisKeyEnum;
import com.njcn.redis.utils.RedisUtil;
import com.njcn.system.api.RocketMqLogFeignClient;
@@ -29,16 +28,57 @@ import java.util.Objects;
* @version V1.0.0
*/
@Component
@RocketMQMessageListener(
topic = "Real_Time_Data_Topic",
consumerGroup = "real_time_consumer",
selectorExpression = "*",
consumeThreadNumber = 10,
enableMsgTrace = true
)
@Slf4j
public class RealTimeDataConsumer {
public class RealTimeDataConsumer extends EnhanceConsumerMessageHandler<MessageDataDTO> implements RocketMQListener<String> {
@Resource
private RedisUtil redisUtil;
@Resource
private RocketMqLogFeignClient rocketMqLogFeignClient;
@Override
public void onMessage(String message) {
MessageDataDTO messageDataDTO = JSONObject.parseObject(message,MessageDataDTO.class);
super.dispatchMessage(messageDataDTO);
}
@MqListener(topic = "Real_Time_Data_Topic", group = "real_time_consumer")
/***
* 通过redis分布式锁判断当前消息所处状态
* 1、null 查不到该key的数据属于第一次消费放行
* 2、fail 上次消息消费时发生异常,放行
* 3、being processed 正在处理,打回去
* 4、success 最近72小时消费成功避免重复消费打回去
*/
@Override
public boolean filter(MessageDataDTO message) {
String keyStatus = redisUtil.getStringByKey(RedisKeyPrefix.REAL_TIME_DATA.concat(message.getKey()));
if (Objects.isNull(keyStatus) || keyStatus.equalsIgnoreCase(MessageStatus.FAIL)) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.REAL_TIME_DATA.concat(message.getKey()), MessageStatus.BEING_PROCESSED, 30L);
return false;
}
return true;
}
/**
* 消费成功缓存到redis5分钟避免重复消费
*/
@Override
protected void consumeSuccess(MessageDataDTO message) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.REAL_TIME_DATA.concat(message.getKey()), MessageStatus.SUCCESS, 5*60L);
}
@Override
protected void handleMessage(MessageDataDTO message) {
String lineId = message.getMonitor();
WebSocketServer.sendInfo(JSONObject.toJSONString(message),lineId);
@@ -49,8 +89,54 @@ public class RealTimeDataConsumer {
/**
* 发生异常时,进行错误信息入库保存
* 默认没有实现类子类可以实现该方法调用feign接口进行入库保存
*/
@Override
protected void saveExceptionMsgLog(MessageDataDTO message, String identity, Exception exception) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.REAL_TIME_DATA.concat(message.getKey()), MessageStatus.FAIL, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
RocketmqMsgErrorLog rocketmqMsgErrorLog = new RocketmqMsgErrorLog();
rocketmqMsgErrorLog.setMsgKey(message.getKey());
rocketmqMsgErrorLog.setResource(message.getSource());
if (identity.equalsIgnoreCase(EnhanceMessageConstant.IDENTITY_SINGLE)) {
//数据库字段配置长度200避免插入失败大致分析异常原因
String exceptionMsg = exception.getMessage();
if(exceptionMsg.length() > 200){
exceptionMsg = exceptionMsg.substring(0,180);
}
rocketmqMsgErrorLog.setRecord(exceptionMsg);
//如果是当前消息重试的则略过
if(!message.getSource().startsWith(EnhanceMessageConstant.RETRY_PREFIX)){
//单次消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
} else {
rocketmqMsgErrorLog.setRecord("重试消费" + super.getMaxRetryTimes() + "次,依旧消费失败。");
//重试N次后依然消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
}
/***
* 处理失败后,是否重试
* 一般开启
*/
@Override
protected boolean isRetry() {
return true;
}
/***
* 消费失败是否抛出异常,抛出异常后就不再消费了
*/
@Override
protected boolean throwException() {
return false;
}

View File

@@ -6,7 +6,6 @@ import com.njcn.message.constant.RedisKeyPrefix;
import com.njcn.message.messagedto.FrontLogslMessage;
import com.njcn.middle.rocket.constant.EnhanceMessageConstant;
import com.njcn.middle.rocket.handler.EnhanceConsumerMessageHandler;
import com.njcn.mq.annotation.MqListener;
import com.njcn.redis.pojo.enums.AppRedisKey;
import com.njcn.redis.pojo.enums.RedisKeyEnum;
import com.njcn.redis.utils.RedisUtil;
@@ -31,15 +30,58 @@ import java.util.Objects;
* @version V1.0.0
*/
@Component
@RocketMQMessageListener(
topic = "log_Topic",
consumerGroup = "Log_Topic_Consumer",
selectorExpression = "*",
consumeThreadNumber = 10,
enableMsgTrace = true
)
@Slf4j
public class TopicLogsConsumer {
public class TopicLogsConsumer extends EnhanceConsumerMessageHandler<FrontLogslMessage> implements RocketMQListener<String> {
@Resource
private RedisUtil redisUtil;
@Resource
private FrontLogsFeignClient frontLogsFeignClient;
@MqListener(topic = "log_Topic", group = "Log_Topic_Consumer")
@Resource
private RocketMqLogFeignClient rocketMqLogFeignClient;
@Override
public void onMessage(String message) {
FrontLogslMessage frontLogslMessage = JSONObject.parseObject(message,FrontLogslMessage.class);
super.dispatchMessage(frontLogslMessage);
}
/***
* 通过redis分布式锁判断当前消息所处状态
* 1、null 查不到该key的数据属于第一次消费放行
* 2、fail 上次消息消费时发生异常,放行
* 3、being processed 正在处理,打回去
* 4、success 最近72小时消费成功避免重复消费打回去
*/
@Override
public boolean filter(FrontLogslMessage message) {
// String keyStatus = redisUtil.getStringByKey(AppRedisKey.RMQ_CONSUME_KEY.concat(message.getKey()));
// if (Objects.isNull(keyStatus) || keyStatus.equalsIgnoreCase(MessageStatus.FAIL)) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.TOPIC_REPLY.concat(message.getKey()), MessageStatus.BEING_PROCESSED, 30L);
// return false;
// }
return false;
}
/**
* 消费成功缓存到redis72小时避免重复消费
*/
@Override
protected void consumeSuccess(FrontLogslMessage message) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.TOPIC_REPLY.concat(message.getKey()), MessageStatus.SUCCESS, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
}
@Override
protected void handleMessage(FrontLogslMessage message) {
//业务处理
@@ -51,8 +93,54 @@ public class TopicLogsConsumer {
/**
* 发生异常时,进行错误信息入库保存
* 默认没有实现类子类可以实现该方法调用feign接口进行入库保存
*/
@Override
protected void saveExceptionMsgLog(FrontLogslMessage message, String identity, Exception exception) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.TOPIC_REPLY.concat(message.getKey()), MessageStatus.FAIL, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
RocketmqMsgErrorLog rocketmqMsgErrorLog = new RocketmqMsgErrorLog();
rocketmqMsgErrorLog.setMsgKey(message.getKey());
rocketmqMsgErrorLog.setResource(message.getSource());
if (identity.equalsIgnoreCase(EnhanceMessageConstant.IDENTITY_SINGLE)) {
//数据库字段配置长度200避免插入失败大致分析异常原因
String exceptionMsg = exception.getMessage();
if(exceptionMsg.length() > 200){
exceptionMsg = exceptionMsg.substring(0,180);
}
rocketmqMsgErrorLog.setRecord(exceptionMsg);
//如果是当前消息重试的则略过
if(!message.getSource().startsWith(EnhanceMessageConstant.RETRY_PREFIX)){
//单次消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
} else {
rocketmqMsgErrorLog.setRecord("重试消费" + super.getMaxRetryTimes() + "次,依旧消费失败。");
//重试N次后依然消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
}
/***
* 处理失败后,是否重试
* 一般开启
*/
@Override
protected boolean isRetry() {
return false;
}
/***
* 消费失败是否抛出异常,抛出异常后就不再消费了
*/
@Override
protected boolean throwException() {
return false;
}

View File

@@ -7,7 +7,6 @@ import com.njcn.message.messagedto.TopicReplyDTO;
import com.njcn.message.constant.RedisKeyPrefix;
import com.njcn.middle.rocket.constant.EnhanceMessageConstant;
import com.njcn.middle.rocket.handler.EnhanceConsumerMessageHandler;
import com.njcn.mq.annotation.MqListener;
import com.njcn.redis.pojo.enums.RedisKeyEnum;
import com.njcn.redis.utils.RedisUtil;
import com.njcn.stat.api.MessAnalysisFeignClient;
@@ -38,14 +37,50 @@ import java.util.Objects;
enableMsgTrace = true
)
@Slf4j
public class TopicReplyConsumer {
public class TopicReplyConsumer extends EnhanceConsumerMessageHandler<TopicReplyDTO> implements RocketMQListener<String> {
@Resource
private RedisUtil redisUtil;
@Resource
private RocketMqLogFeignClient rocketMqLogFeignClient;
@Autowired
private MessAnalysisFeignClient messAnalysisFeignClient;
@MqListener(topic = "Topic_Reply_Topic", group = "Topic_Reply_Topic_Consumer")
@Override
public void onMessage(String message) {
TopicReplyDTO topicReplyDTO = JSONObject.parseObject(message,TopicReplyDTO.class);
super.dispatchMessage(topicReplyDTO);
}
/***
* 通过redis分布式锁判断当前消息所处状态
* 1、null 查不到该key的数据属于第一次消费放行
* 2、fail 上次消息消费时发生异常,放行
* 3、being processed 正在处理,打回去
* 4、success 最近72小时消费成功避免重复消费打回去
*/
@Override
public boolean filter(TopicReplyDTO message) {
// String keyStatus = redisUtil.getStringByKey(AppRedisKey.RMQ_CONSUME_KEY.concat(message.getKey()));
// if (Objects.isNull(keyStatus) || keyStatus.equalsIgnoreCase(MessageStatus.FAIL)) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.TOPIC_REPLY.concat(message.getKey()), MessageStatus.BEING_PROCESSED, 30L);
// return false;
// }
return false;
}
/**
* 消费成功缓存到redis72小时避免重复消费
*/
@Override
protected void consumeSuccess(TopicReplyDTO message) {
// redisUtil.saveByKeyWithExpire(RedisKeyPrefix.TOPIC_REPLY.concat(message.getKey()), MessageStatus.SUCCESS, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
}
@Override
protected void handleMessage(TopicReplyDTO message) {
//“12345”补招回复
if(Objects.equals(message.getGuid(),"12345")){
@@ -59,4 +94,56 @@ public class TopicReplyConsumer {
/**
* 发生异常时,进行错误信息入库保存
* 默认没有实现类子类可以实现该方法调用feign接口进行入库保存
*/
@Override
protected void saveExceptionMsgLog(TopicReplyDTO message, String identity, Exception exception) {
redisUtil.saveByKeyWithExpire(RedisKeyPrefix.TOPIC_REPLY.concat(message.getKey()), MessageStatus.FAIL, RedisKeyEnum.ROCKET_MQ_KEY.getTime());
RocketmqMsgErrorLog rocketmqMsgErrorLog = new RocketmqMsgErrorLog();
rocketmqMsgErrorLog.setMsgKey(message.getKey());
rocketmqMsgErrorLog.setResource(message.getSource());
if (identity.equalsIgnoreCase(EnhanceMessageConstant.IDENTITY_SINGLE)) {
//数据库字段配置长度200避免插入失败大致分析异常原因
String exceptionMsg = exception.getMessage();
if(exceptionMsg.length() > 200){
exceptionMsg = exceptionMsg.substring(0,180);
}
rocketmqMsgErrorLog.setRecord(exceptionMsg);
//如果是当前消息重试的则略过
if(!message.getSource().startsWith(EnhanceMessageConstant.RETRY_PREFIX)){
//单次消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
} else {
rocketmqMsgErrorLog.setRecord("重试消费" + super.getMaxRetryTimes() + "次,依旧消费失败。");
//重试N次后依然消费异常
rocketMqLogFeignClient.add(rocketmqMsgErrorLog);
}
}
/***
* 处理失败后,是否重试
* 一般开启
*/
@Override
protected boolean isRetry() {
return true;
}
/***
* 消费失败是否抛出异常,抛出异常后就不再消费了
*/
@Override
protected boolean throwException() {
return false;
}
}

View File

@@ -1,23 +1,23 @@
//package com.njcn.message.mq;
//
//import com.njcn.message.messagedto.MessageDataDTO;
//import com.njcn.mq.annotation.MqListener;
//import com.njcn.stat.api.MessAnalysisFeignClient;
//import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
//import org.springframework.stereotype.Component;
//
//import javax.annotation.Resource;
//import java.util.Collections;
//
//@Component
//@ConditionalOnProperty(name = "mq.type", havingValue = "redis-stream")
//public class FrontDataMqListener {
//
// @Resource
// private MessAnalysisFeignClient messAnalysisFeignClient;
//
// @MqListener(topic = "LN_Topic", group = "ln_consumer")
// public void onFrontData(MessageDataDTO msg) {
// messAnalysisFeignClient.analysis(Collections.singletonList(msg));
// }
//}
package com.njcn.message.mq;
import com.njcn.message.messagedto.MessageDataDTO;
import com.njcn.mq.annotation.MqListener;
import com.njcn.stat.api.MessAnalysisFeignClient;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.Collections;
@Component
@ConditionalOnProperty(name = "mq.type", havingValue = "redis-stream")
public class FrontDataMqListener {
@Resource
private MessAnalysisFeignClient messAnalysisFeignClient;
@MqListener(topic = "LN_Topic", group = "ln_consumer")
public void onFrontData(MessageDataDTO msg) {
messAnalysisFeignClient.analysis(Collections.singletonList(msg));
}
}

View File

@@ -1,19 +1,19 @@
//package com.njcn.message.mq;
//
//import com.njcn.message.constant.BusinessTopic;
//import com.njcn.message.message.RecallMessage;
//import com.njcn.mq.core.MqTemplate;
//import org.springframework.stereotype.Component;
//
//import javax.annotation.Resource;
//
//@Component
//public class RecallMqSender {
//
// @Resource
// private MqTemplate mqTemplate;
//
// public void send(RecallMessage message, String nodeId) {
// mqTemplate.send(nodeId + "_" + BusinessTopic.RECALL_TOPIC, message);
// }
//}
package com.njcn.message.mq;
import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.RecallMessage;
import com.njcn.mq.core.MqTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
@Component
public class RecallMqSender {
@Resource
private MqTemplate mqTemplate;
public void send(RecallMessage message, String nodeId) {
mqTemplate.send(nodeId + "_" + BusinessTopic.RECALL_TOPIC, message);
}
}

View File

@@ -86,7 +86,7 @@ public class ProduceController extends BaseController {
@ApiImplicitParam(name = "message", value = "参数", required = true)
public HttpResult<String> askFileSys(@RequestBody AskFileSysMessage message){
String methodDescribe = getMethodDescribe("askFileSys");
askFileSysMessaggeTemplate.sendMember(message);
SendResult sendResult = askFileSysMessaggeTemplate.sendMember(message);
return HttpResultUtil.assembleCommonResponseResult(CommonResponseEnum.SUCCESS,message.getGuid() , methodDescribe);
}
}

View File

@@ -7,13 +7,10 @@ import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.AskFileSysMessage;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
/**
* Description:
* Date: 2024/12/13 15:15【需求编号】
@@ -22,14 +19,16 @@ import javax.annotation.Resource;
* @version V1.0.0
*/
@Component
public class AskFileSysMessaggeTemplate {
@Resource
private MqTemplate mqTemplate;
public void sendMember(AskFileSysMessage askRealDataMessage) {
public class AskFileSysMessaggeTemplate extends RocketMQEnhanceTemplate {
public AskFileSysMessaggeTemplate(RocketMQTemplate template) {
super(template);
}
public SendResult sendMember(AskFileSysMessage askRealDataMessage) {
BaseMessage baseMessage = new BaseMessage();
baseMessage.setMessageBody(JSONObject.toJSONString(askRealDataMessage));
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
baseMessage.setKey(askRealDataMessage.getGuid());
mqTemplate.send(BusinessTopic.FILE_TOPIC,askRealDataMessage.getNodeId() , baseMessage);
return send(BusinessTopic.FILE_TOPIC,askRealDataMessage.getNodeId() , baseMessage);
}
}

View File

@@ -7,13 +7,10 @@ import com.njcn.message.constant.BusinessTopic;
import com.njcn.message.message.AskRealDataMessage;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
/**
* Description:
* Date: 2024/12/13 15:15【需求编号】
@@ -22,13 +19,15 @@ import javax.annotation.Resource;
* @version V1.0.0
*/
@Component
public class AskRealDataMessaggeTemplate {
@Resource
private MqTemplate mqTemplate;
public void sendMember(BaseMessage askRealDataMessage,String nodeId) {
public class AskRealDataMessaggeTemplate extends RocketMQEnhanceTemplate {
public AskRealDataMessaggeTemplate(RocketMQTemplate template) {
super(template);
}
public SendResult sendMember(BaseMessage askRealDataMessage,String nodeId) {
askRealDataMessage.setSource(BusinessResource.WEB_RESOURCE);
AskRealDataMessage dto = JSON.parseObject(askRealDataMessage.getMessageBody(), AskRealDataMessage.class);
askRealDataMessage.setKey(dto.getLine());
mqTemplate.send(BusinessTopic.ASK_REAL_DATA_TOPIC,nodeId , askRealDataMessage);
return send(BusinessTopic.ASK_REAL_DATA_TOPIC,nodeId , askRealDataMessage);
}
}

View File

@@ -7,13 +7,10 @@ import com.njcn.message.message.AskRealDataMessage;
import com.njcn.message.message.DeviceRebootMessage;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
/**
* 类的介绍:
*
@@ -22,12 +19,14 @@ import javax.annotation.Resource;
* @createTime 2023/8/11 15:28
*/
@Component
public class DeviceRebootMessageTemplate {
public class DeviceRebootMessageTemplate extends RocketMQEnhanceTemplate {
@Resource
private MqTemplate mqTemplate;
public void sendMember(BaseMessage baseMessage,String nodeId) {
public DeviceRebootMessageTemplate(RocketMQTemplate template) {
super(template);
}
public SendResult sendMember(BaseMessage baseMessage,String nodeId) {
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
mqTemplate.send(BusinessTopic.CONTROL_TOPIC,nodeId, baseMessage);
return send(BusinessTopic.CONTROL_TOPIC,nodeId, baseMessage);
}
}

View File

@@ -7,13 +7,10 @@ import com.njcn.message.message.AskRealDataMessage;
import com.njcn.message.message.ProcessRebootMessage;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
/**
* 类的介绍:
*
@@ -22,14 +19,16 @@ import javax.annotation.Resource;
* @createTime 2023/8/11 15:28
*/
@Component
public class ProcessRebootMessageTemplate {
public class ProcessRebootMessageTemplate extends RocketMQEnhanceTemplate {
@Resource
private MqTemplate mqTemplate;
public void sendMember(BaseMessage baseMessage,String nodeId) {
public ProcessRebootMessageTemplate(RocketMQTemplate template) {
super(template);
}
public SendResult sendMember(BaseMessage baseMessage,String nodeId) {
baseMessage.setSource(BusinessResource.WEB_RESOURCE);
ProcessRebootMessage dto = JSON.parseObject(baseMessage.getMessageBody(), ProcessRebootMessage.class);
baseMessage.setKey(dto.getIndex()+"");
mqTemplate.send(BusinessTopic.PROCESS_TOPIC,nodeId, baseMessage);
return send(BusinessTopic.PROCESS_TOPIC,nodeId, baseMessage);
}
}

View File

@@ -5,13 +5,10 @@ import com.njcn.message.constant.BusinessResource;
import com.njcn.message.constant.BusinessTopic;
import com.njcn.middle.rocket.domain.BaseMessage;
import com.njcn.middle.rocket.template.RocketMQEnhanceTemplate;
import com.njcn.mq.core.MqTemplate;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
/**
* Description:
* Date: 2024/12/13 15:15【需求编号】
@@ -20,11 +17,13 @@ import javax.annotation.Resource;
* @version V1.0.0
*/
@Component
public class RecallMessaggeTemplate {
@Resource
private MqTemplate mqTemplate;
public void sendMember(BaseMessage recallMessage,String nodeId) {
public class RecallMessaggeTemplate extends RocketMQEnhanceTemplate {
public RecallMessaggeTemplate(RocketMQTemplate template) {
super(template);
}
public SendResult sendMember(BaseMessage recallMessage,String nodeId) {
recallMessage.setSource(BusinessResource.WEB_RESOURCE);
mqTemplate.send(BusinessTopic.RECALL_TOPIC,nodeId , recallMessage);
return send(BusinessTopic.RECALL_TOPIC,nodeId , recallMessage);
}
}

View File

@@ -1,5 +1,3 @@
spring:
profiles:
active: @spring.profiles.active@
mq:
type: redis-stream
active: @spring.profiles.active@

View File

@@ -1,210 +1,210 @@
//package com.njcn.message.mq;
//
//import com.alibaba.fastjson.JSON;
//import com.njcn.message.message.RecallMessage;
//import com.njcn.message.messagedto.MessageDataDTO;
//import com.njcn.middle.stream.autoconfig.RedisStreamAutoConfiguration;
//import com.njcn.mq.autoconfig.MqCoreAutoConfiguration;
//import com.njcn.mq.container.MqListenerRegistry;
//import com.njcn.mq.core.MqTemplate;
//import com.njcn.mq.driver.redis.RedisStreamMqDriver;
//import com.njcn.mq.driver.redis.RedisStreamMqDriverAutoConfiguration;
//import com.njcn.mq.driver.rocketmq.RocketMqDriver;
//import com.njcn.mq.driver.rocketmq.RocketMqDriverAutoConfiguration;
//import com.njcn.mq.spi.MqDriver;
//import com.njcn.stat.api.MessAnalysisFeignClient;
//import org.apache.rocketmq.spring.autoconfigure.RocketMQProperties;
//import org.apache.rocketmq.spring.core.RocketMQTemplate;
//import org.junit.jupiter.api.Assumptions;
//import org.junit.jupiter.api.Test;
//import org.springframework.boot.autoconfigure.AutoConfigurations;
//import org.springframework.boot.autoconfigure.data.redis.RedisAutoConfiguration;
//import org.springframework.boot.test.context.runner.ApplicationContextRunner;
//import org.springframework.data.redis.connection.stream.MapRecord;
//import org.springframework.data.redis.connection.stream.ReadOffset;
//import org.springframework.data.redis.connection.stream.StreamOffset;
//import org.springframework.data.redis.connection.stream.StreamReadOptions;
//import org.springframework.data.redis.connection.stream.StreamRecords;
//import org.springframework.data.redis.core.StringRedisTemplate;
//
//import java.net.InetSocketAddress;
//import java.net.Socket;
//import java.util.HashMap;
//import java.util.List;
//import java.util.Map;
//
//import static org.assertj.core.api.Assertions.assertThat;
//import static org.mockito.ArgumentMatchers.anyList;
//import static org.mockito.Mockito.mock;
//import static org.mockito.Mockito.timeout;
//import static org.mockito.Mockito.verify;
//
///**
// * MQ starter 切片接入验证测试ApplicationContextRunner不启动全应用
// * 真实 Redis 要求localhost:6379docker redis-stream-it仅上/下行收发用例需要,
// * 纯装配用例不需要Lettuce 懒连接,不触发实际连接)。
// *
// * <p>装配一律用 {@code withConfiguration(AutoConfigurations.of(...))},让 Spring 按
// * <b>真实 autoconfig 顺序</b>AutoConfigurationSorter处理与生产环境一致。以此
// * 验证 {@link MqCoreAutoConfiguration} 的 {@code @ConditionalOnBean(MqDriver)} 能在
// * 真实顺序下拿到 driver bean —— 这依赖 driver autoconfig 上的
// * {@code @AutoConfigureBefore(MqCoreAutoConfiguration.class)}。若缺该保障,按类名字母序
// * {@code com.njcn.mq.autoconfig.*} 先于 {@code com.njcn.mq.driver.*}core 会在 driver
// * bean 注册前求值条件失败,导致 MqTemplate / MqListenerRegistry 不装配。
// */
//public class MqSliceWiringTest {
//
// private static boolean redisAvailable() {
// try {
// Socket s = new Socket();
// s.connect(new InetSocketAddress("localhost", 6379), 500);
// s.close();
// return true;
// } catch (Exception e) {
// return false;
// }
// }
//
// /** redis-stream 模式:真实 autoconfig 顺序装 5 个 autoconfig + 业务 bean上/下行用)。 */
// private ApplicationContextRunner redisStreamRunner() {
// return new ApplicationContextRunner()
// .withConfiguration(AutoConfigurations.of(
// RedisAutoConfiguration.class,
// RedisStreamAutoConfiguration.class,
// RedisStreamMqDriverAutoConfiguration.class,
// RocketMqDriverAutoConfiguration.class,
// MqCoreAutoConfiguration.class))
// .withUserConfiguration(FrontDataMqListener.class, RecallMqSender.class)
// .withBean(MessAnalysisFeignClient.class,
// () -> mock(MessAnalysisFeignClient.class))
// .withPropertyValues(
// "mq.type=redis-stream",
// "spring.redis.host=localhost",
// "spring.redis.port=6379");
// }
//
// // ------------------------------------------------------------------
// // 1. 装配mq.type 二选一 + core bean 在真实 autoconfig 顺序下装配
// // ------------------------------------------------------------------
//
// @Test
// void testRedisStreamDriverSelected() {
// // 纯装配检查,不依赖真实 RedisLettuce 懒连接)。
// new ApplicationContextRunner()
// .withConfiguration(AutoConfigurations.of(
// RedisAutoConfiguration.class,
// RedisStreamAutoConfiguration.class,
// RedisStreamMqDriverAutoConfiguration.class,
// RocketMqDriverAutoConfiguration.class,
// MqCoreAutoConfiguration.class))
// .withPropertyValues(
// "mq.type=redis-stream",
// "spring.redis.host=localhost",
// "spring.redis.port=6379")
// .run(ctx -> {
// assertThat(ctx).hasSingleBean(MqDriver.class);
// assertThat(ctx.getBean(MqDriver.class)).isInstanceOf(RedisStreamMqDriver.class);
// // 真实顺序下 core 的 @ConditionalOnBean(MqDriver) 必须命中
// assertThat(ctx).hasSingleBean(MqTemplate.class);
// assertThat(ctx).hasSingleBean(MqListenerRegistry.class);
// });
// }
//
// @Test
// void testRocketMqDriverSelected() {
// // rocket 测试不依赖真实 Redis不装 RedisAutoConfigurationredis driver 因 mq.type 不匹配而 inactive。
// new ApplicationContextRunner()
// .withConfiguration(AutoConfigurations.of(
// RedisStreamMqDriverAutoConfiguration.class,
// RocketMqDriverAutoConfiguration.class,
// MqCoreAutoConfiguration.class))
// .withBean(RocketMQTemplate.class, () -> mock(RocketMQTemplate.class))
// .withBean(RocketMQProperties.class, RocketMQProperties::new)
// .withPropertyValues("mq.type=rocketmq")
// .run(ctx -> {
// assertThat(ctx).hasSingleBean(MqDriver.class);
// assertThat(ctx.getBean(MqDriver.class)).isInstanceOf(RocketMqDriver.class);
// assertThat(ctx).hasSingleBean(MqTemplate.class);
// assertThat(ctx).hasSingleBean(MqListenerRegistry.class);
// });
// }
//
// // ------------------------------------------------------------------
// // 2. 上行(需要真实 Redis
// // 向 LN_Topic stream XADD 一条消息,验证 FrontDataMqListener 触发了
// // messAnalysisFeignClient.analysis(...)
// // ------------------------------------------------------------------
//
// @Test
// void testUpstreamFrontData() {
// Assumptions.assumeTrue(redisAvailable(), "Redis not available at localhost:6379, skipping");
// redisStreamRunner().run(ctx -> {
// StringRedisTemplate redis = ctx.getBean(StringRedisTemplate.class);
// MessAnalysisFeignClient mockClient = ctx.getBean(MessAnalysisFeignClient.class);
//
// // 构造 RedisStreamMqDriver 消费时期望的字段格式
// // enc=plain → body 即为原始 JSON参见 MessageCodec.encodeBody(json, false)
// // stream key = StreamKeyBuilder.buildStream("LN_Topic") = "LN_Topic"
// // (默认 envIsolation=false无 env 前缀)
// MessageDataDTO dto = new MessageDataDTO();
// dto.setDataType(1);
// String json = JSON.toJSONString(dto);
//
// Map<String, String> fields = new HashMap<String, String>();
// fields.put("key", "test-key-upstream-1");
// fields.put("tag", "");
// fields.put("enc", "plain");
// fields.put("body", json);
// fields.put("ts", String.valueOf(System.currentTimeMillis()));
//
// // XADD 到 LN_TopicMqListenerRegistry.start() 已在 context 启动时订阅该 stream
// redis.opsForStream().add(
// StreamRecords.newRecord().in("LN_Topic").ofMap(fields)
// );
//
// // 等待消费线程处理blockMs 默认 2000mstimeout=5000ms 足够)
// verify(mockClient, timeout(5000)).analysis(anyList());
// });
// }
//
// // ------------------------------------------------------------------
// // 3. 下行(需要真实 Redis
// // 通过 RecallMqSender 发送 RecallMessage验证 stream 内容正确
// // ------------------------------------------------------------------
//
// @Test
// void testDownstreamRecall() {
// Assumptions.assumeTrue(redisAvailable(), "Redis not available at localhost:6379, skipping");
// redisStreamRunner().run(ctx -> {
// StringRedisTemplate redis = ctx.getBean(StringRedisTemplate.class);
// RecallMqSender sender = ctx.getBean(RecallMqSender.class);
//
// // 清理,确保只有本次发送的消息
// redis.delete("n1_recall_Topic");
//
// RecallMessage msg = new RecallMessage();
// msg.setGuid("g1");
// // MqTemplate.send("n1_recall_Topic", msg) → RedisStreamMqDriver XADD
// // stream key = StreamKeyBuilder.buildStream("n1_recall_Topic") = "n1_recall_Topic"
// sender.send(msg, "n1");
//
// Long size = redis.opsForStream().size("n1_recall_Topic");
// assertThat(size).isGreaterThanOrEqualTo(1L);
//
// @SuppressWarnings("unchecked")
// List<MapRecord<String, Object, Object>> records =
// redis.opsForStream().read(
// StreamReadOptions.empty().count(1L),
// StreamOffset.create("n1_recall_Topic", ReadOffset.from("0-0"))
// );
// assertThat(records).isNotEmpty();
//
// Map<Object, Object> firstFields = records.get(0).getValue();
// // DefaultMqTemplate.send(topic, payload) 使用 encodeBody(json, false) → enc="plain"
// // 因此 body 即为原始 JSON无需 gzip 解压
// String body = (String) firstFields.get("body");
// RecallMessage result = JSON.parseObject(body, RecallMessage.class);
// assertThat(result.getGuid()).isEqualTo("g1");
// });
// }
//}
package com.njcn.message.mq;
import com.alibaba.fastjson.JSON;
import com.njcn.message.message.RecallMessage;
import com.njcn.message.messagedto.MessageDataDTO;
import com.njcn.middle.stream.autoconfig.RedisStreamAutoConfiguration;
import com.njcn.mq.autoconfig.MqCoreAutoConfiguration;
import com.njcn.mq.container.MqListenerRegistry;
import com.njcn.mq.core.MqTemplate;
import com.njcn.mq.driver.redis.RedisStreamMqDriver;
import com.njcn.mq.driver.redis.RedisStreamMqDriverAutoConfiguration;
import com.njcn.mq.driver.rocketmq.RocketMqDriver;
import com.njcn.mq.driver.rocketmq.RocketMqDriverAutoConfiguration;
import com.njcn.mq.spi.MqDriver;
import com.njcn.stat.api.MessAnalysisFeignClient;
import org.apache.rocketmq.spring.autoconfigure.RocketMQProperties;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.junit.jupiter.api.Assumptions;
import org.junit.jupiter.api.Test;
import org.springframework.boot.autoconfigure.AutoConfigurations;
import org.springframework.boot.autoconfigure.data.redis.RedisAutoConfiguration;
import org.springframework.boot.test.context.runner.ApplicationContextRunner;
import org.springframework.data.redis.connection.stream.MapRecord;
import org.springframework.data.redis.connection.stream.ReadOffset;
import org.springframework.data.redis.connection.stream.StreamOffset;
import org.springframework.data.redis.connection.stream.StreamReadOptions;
import org.springframework.data.redis.connection.stream.StreamRecords;
import org.springframework.data.redis.core.StringRedisTemplate;
import java.net.InetSocketAddress;
import java.net.Socket;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
/**
* MQ starter 切片接入验证测试ApplicationContextRunner不启动全应用
* 真实 Redis 要求localhost:6379docker redis-stream-it仅上/下行收发用例需要,
* 纯装配用例不需要Lettuce 懒连接,不触发实际连接)。
*
* <p>装配一律用 {@code withConfiguration(AutoConfigurations.of(...))},让 Spring 按
* <b>真实 autoconfig 顺序</b>AutoConfigurationSorter处理与生产环境一致。以此
* 验证 {@link MqCoreAutoConfiguration} 的 {@code @ConditionalOnBean(MqDriver)} 能在
* 真实顺序下拿到 driver bean —— 这依赖 driver autoconfig 上的
* {@code @AutoConfigureBefore(MqCoreAutoConfiguration.class)}。若缺该保障,按类名字母序
* {@code com.njcn.mq.autoconfig.*} 先于 {@code com.njcn.mq.driver.*}core 会在 driver
* bean 注册前求值条件失败,导致 MqTemplate / MqListenerRegistry 不装配。
*/
public class MqSliceWiringTest {
private static boolean redisAvailable() {
try {
Socket s = new Socket();
s.connect(new InetSocketAddress("localhost", 6379), 500);
s.close();
return true;
} catch (Exception e) {
return false;
}
}
/** redis-stream 模式:真实 autoconfig 顺序装 5 个 autoconfig + 业务 bean上/下行用)。 */
private ApplicationContextRunner redisStreamRunner() {
return new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(
RedisAutoConfiguration.class,
RedisStreamAutoConfiguration.class,
RedisStreamMqDriverAutoConfiguration.class,
RocketMqDriverAutoConfiguration.class,
MqCoreAutoConfiguration.class))
.withUserConfiguration(FrontDataMqListener.class, RecallMqSender.class)
.withBean(MessAnalysisFeignClient.class,
() -> mock(MessAnalysisFeignClient.class))
.withPropertyValues(
"mq.type=redis-stream",
"spring.redis.host=localhost",
"spring.redis.port=6379");
}
// ------------------------------------------------------------------
// 1. 装配mq.type 二选一 + core bean 在真实 autoconfig 顺序下装配
// ------------------------------------------------------------------
@Test
void testRedisStreamDriverSelected() {
// 纯装配检查,不依赖真实 RedisLettuce 懒连接)。
new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(
RedisAutoConfiguration.class,
RedisStreamAutoConfiguration.class,
RedisStreamMqDriverAutoConfiguration.class,
RocketMqDriverAutoConfiguration.class,
MqCoreAutoConfiguration.class))
.withPropertyValues(
"mq.type=redis-stream",
"spring.redis.host=localhost",
"spring.redis.port=6379")
.run(ctx -> {
assertThat(ctx).hasSingleBean(MqDriver.class);
assertThat(ctx.getBean(MqDriver.class)).isInstanceOf(RedisStreamMqDriver.class);
// 真实顺序下 core 的 @ConditionalOnBean(MqDriver) 必须命中
assertThat(ctx).hasSingleBean(MqTemplate.class);
assertThat(ctx).hasSingleBean(MqListenerRegistry.class);
});
}
@Test
void testRocketMqDriverSelected() {
// rocket 测试不依赖真实 Redis不装 RedisAutoConfigurationredis driver 因 mq.type 不匹配而 inactive。
new ApplicationContextRunner()
.withConfiguration(AutoConfigurations.of(
RedisStreamMqDriverAutoConfiguration.class,
RocketMqDriverAutoConfiguration.class,
MqCoreAutoConfiguration.class))
.withBean(RocketMQTemplate.class, () -> mock(RocketMQTemplate.class))
.withBean(RocketMQProperties.class, RocketMQProperties::new)
.withPropertyValues("mq.type=rocketmq")
.run(ctx -> {
assertThat(ctx).hasSingleBean(MqDriver.class);
assertThat(ctx.getBean(MqDriver.class)).isInstanceOf(RocketMqDriver.class);
assertThat(ctx).hasSingleBean(MqTemplate.class);
assertThat(ctx).hasSingleBean(MqListenerRegistry.class);
});
}
// ------------------------------------------------------------------
// 2. 上行(需要真实 Redis
// 向 LN_Topic stream XADD 一条消息,验证 FrontDataMqListener 触发了
// messAnalysisFeignClient.analysis(...)
// ------------------------------------------------------------------
@Test
void testUpstreamFrontData() {
Assumptions.assumeTrue(redisAvailable(), "Redis not available at localhost:6379, skipping");
redisStreamRunner().run(ctx -> {
StringRedisTemplate redis = ctx.getBean(StringRedisTemplate.class);
MessAnalysisFeignClient mockClient = ctx.getBean(MessAnalysisFeignClient.class);
// 构造 RedisStreamMqDriver 消费时期望的字段格式
// enc=plain → body 即为原始 JSON参见 MessageCodec.encodeBody(json, false)
// stream key = StreamKeyBuilder.buildStream("LN_Topic") = "LN_Topic"
// (默认 envIsolation=false无 env 前缀)
MessageDataDTO dto = new MessageDataDTO();
dto.setDataType(1);
String json = JSON.toJSONString(dto);
Map<String, String> fields = new HashMap<String, String>();
fields.put("key", "test-key-upstream-1");
fields.put("tag", "");
fields.put("enc", "plain");
fields.put("body", json);
fields.put("ts", String.valueOf(System.currentTimeMillis()));
// XADD 到 LN_TopicMqListenerRegistry.start() 已在 context 启动时订阅该 stream
redis.opsForStream().add(
StreamRecords.newRecord().in("LN_Topic").ofMap(fields)
);
// 等待消费线程处理blockMs 默认 2000mstimeout=5000ms 足够)
verify(mockClient, timeout(5000)).analysis(anyList());
});
}
// ------------------------------------------------------------------
// 3. 下行(需要真实 Redis
// 通过 RecallMqSender 发送 RecallMessage验证 stream 内容正确
// ------------------------------------------------------------------
@Test
void testDownstreamRecall() {
Assumptions.assumeTrue(redisAvailable(), "Redis not available at localhost:6379, skipping");
redisStreamRunner().run(ctx -> {
StringRedisTemplate redis = ctx.getBean(StringRedisTemplate.class);
RecallMqSender sender = ctx.getBean(RecallMqSender.class);
// 清理,确保只有本次发送的消息
redis.delete("n1_recall_Topic");
RecallMessage msg = new RecallMessage();
msg.setGuid("g1");
// MqTemplate.send("n1_recall_Topic", msg) → RedisStreamMqDriver XADD
// stream key = StreamKeyBuilder.buildStream("n1_recall_Topic") = "n1_recall_Topic"
sender.send(msg, "n1");
Long size = redis.opsForStream().size("n1_recall_Topic");
assertThat(size).isGreaterThanOrEqualTo(1L);
@SuppressWarnings("unchecked")
List<MapRecord<String, Object, Object>> records =
redis.opsForStream().read(
StreamReadOptions.empty().count(1L),
StreamOffset.create("n1_recall_Topic", ReadOffset.from("0-0"))
);
assertThat(records).isNotEmpty();
Map<Object, Object> firstFields = records.get(0).getValue();
// DefaultMqTemplate.send(topic, payload) 使用 encodeBody(json, false) → enc="plain"
// 因此 body 即为原始 JSON无需 gzip 解压
String body = (String) firstFields.get("body");
RecallMessage result = JSON.parseObject(body, RecallMessage.class);
assertThat(result.getGuid()).isEqualTo("g1");
});
}
}

View File

@@ -34,8 +34,8 @@
<properties>
<spring.profiles.active>sjzx</spring.profiles.active>
<!--内网-->
<middle.server.url>192.168.1.68</middle.server.url>
<service.server.url>192.168.2.130</service.server.url>
<middle.server.url>192.168.1.65</middle.server.url>
<service.server.url>192.168.1.65</service.server.url>
<docker.server.url>192.168.1.22</docker.server.url>
<nacos.url>${middle.server.url}:18848</nacos.url>
<nacos.namespace></nacos.namespace>