From 3de93e8d66f0de3b4ba1ba74f5801816daa91c7c Mon Sep 17 00:00:00 2001 From: bgy Date: Mon, 6 Jul 2026 15:40:47 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96MQ=E7=9B=91=E6=8E=A7=E5=8A=9F?= =?UTF-8?q?=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../scheduletask/mapper/TopsailTransmitLogDao.java | 12 ++-- .../topsail/scheduletask/pojo/FocusMqCount.java | 35 ++++++++++ .../scheduletask/pojo/TopsailTransmitLog.java | 4 ++ .../task/CheckRabbitMqScheduleTask.java | 78 ++++++++++++++++------ .../mapper/TopsailTransmitLogMapper.xml | 15 ++++- 5 files changed, 118 insertions(+), 26 deletions(-) create mode 100644 src/main/java/com/topsail/scheduletask/pojo/FocusMqCount.java diff --git a/src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java b/src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java index 514788a..49e67da 100644 --- a/src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java +++ b/src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java @@ -1,9 +1,6 @@ package com.topsail.scheduletask.mapper; -import com.topsail.scheduletask.pojo.DeviceOnlineInfo; -import com.topsail.scheduletask.pojo.FocusDeviceOnline; -import com.topsail.scheduletask.pojo.TopsailProductData; -import com.topsail.scheduletask.pojo.TopsailTransmitLog; +import com.topsail.scheduletask.pojo.*; import org.apache.ibatis.annotations.Param; import org.springframework.stereotype.Repository; @@ -79,4 +76,11 @@ public interface TopsailTransmitLogDao { * @return 影响行数 */ List getProductDataDevice(@Param("isProductData") Integer isProductData); + + /** + * 查询关注的队列消息数情况 + * @param isFocus + * @return + */ + List getIgnoreQueue(@Param("isFocus") Integer isFocus); } diff --git a/src/main/java/com/topsail/scheduletask/pojo/FocusMqCount.java b/src/main/java/com/topsail/scheduletask/pojo/FocusMqCount.java new file mode 100644 index 0000000..bcee55b --- /dev/null +++ b/src/main/java/com/topsail/scheduletask/pojo/FocusMqCount.java @@ -0,0 +1,35 @@ +package com.topsail.scheduletask.pojo; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + + +/** + * 通用明文数据转发对象 + */ +@Data +@AllArgsConstructor +@NoArgsConstructor +@Builder +public class FocusMqCount { + /** + * 队列名称 + */ + private String queuesName; + /** + * 排除的队列名称 + */ + private String excludeName; + + /** + * 是否关注1关注0不关注 + */ + private Integer isFocus; + /** + * 告警接收的邮箱 + */ + private String sendMails; + +} diff --git a/src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java b/src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java index 000f14f..0176058 100644 --- a/src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java +++ b/src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java @@ -102,5 +102,9 @@ public class TopsailTransmitLog { * 上一次转发时间 */ private String lastForwardTime; + /** + * 转发源地址 + */ + private String topForwardFrom; } diff --git a/src/main/java/com/topsail/scheduletask/task/CheckRabbitMqScheduleTask.java b/src/main/java/com/topsail/scheduletask/task/CheckRabbitMqScheduleTask.java index 166727f..04ad53e 100644 --- a/src/main/java/com/topsail/scheduletask/task/CheckRabbitMqScheduleTask.java +++ b/src/main/java/com/topsail/scheduletask/task/CheckRabbitMqScheduleTask.java @@ -5,6 +5,7 @@ import com.topsail.scheduletask.mapper.InformLogDao; import com.topsail.scheduletask.mapper.TopsailTransmitLogDao; import com.topsail.scheduletask.pojo.DeviceOnlineInfo; import com.topsail.scheduletask.pojo.FocusDeviceOnline; +import com.topsail.scheduletask.pojo.FocusMqCount; import com.topsail.scheduletask.service.AmqpService; import com.topsail.scheduletask.util.MaiSenderlUtil; import org.springframework.beans.factory.annotation.Autowired; @@ -42,25 +43,62 @@ public class CheckRabbitMqScheduleTask { //或直接指定时间间隔,例如:5秒 // @Scheduled(fixedRate=5000) private void checkRabbitMqCount() { - //查询rabbitmq的队列名称 - List allQueueNames = amqpService.getAllQueueNames(); - Map queueMessageCountMap = new HashMap<>(); - for (String queueName : allQueueNames) { - //跳过忽略的队列 - if (queueName.equals(ignoreQueue)) { - continue; - } - //获取队列中待消费的消息数量 - long count = amqpService.getQueueMessageCount(queueName); - if (count > 2000) { - System.err.println("队列名称:" + queueName + ",待消费消息数量:" + count); - queueMessageCountMap.put(queueName, count); + try { + //查询忽略的队列名称 + List list = topsailTransmitLogDao.getIgnoreQueue(1); + if (list != null && list.size() > 0) { + for (FocusMqCount focusMqCount : list) { + StringBuilder ignoreQueueBuilder = new StringBuilder(); + ignoreQueueBuilder.append(focusMqCount.getExcludeName()); + //查询rabbitmq的队列名称 + List allQueueNames = amqpService.getAllQueueNames(); + Map queueMessageCountMap = new HashMap<>(); + for (String queueName : allQueueNames) { + //跳过忽略的队列 + if (ignoreQueueBuilder.toString().contains(queueName)) { + continue; + } + //获取队列中待消费的消息数量 + long count = amqpService.getQueueMessageCount(queueName); + if (count > 2000) { + System.err.println("队列名称:" + queueName + ",待消费消息数量:" + count); + queueMessageCountMap.put(queueName, count); + } + } + if (ignoreQueueBuilder.length() > 0) { + ignoreQueueBuilder.delete(0, ignoreQueueBuilder.length()); + } + if (queueMessageCountMap.size() > 0) { + //将queueMessageCountMap转成json字符串 + String json = JSON.toJSONString(queueMessageCountMap); + String subject = "MQ中设备数据流转异常"; + String today = new SimpleDateFormat("yyyy-MM-dd HH:00:00").format(new Date()); + Integer logCount = informLogDao.getNewInformLogCount(subject, today); + if (focusMqCount.getSendMails() != null && logCount != null && logCount < 1) { + if (focusMqCount.getSendMails().contains(",")) { + String[] mails = focusMqCount.getSendMails().split(","); + for (String mail : mails) { + //验证邮箱格式 + if (!maiSenderlUtil.isEmail(mail)) { + continue; + } + //3.发送邮件 + String content = subject + ",待消费消息数量:" + json; + maiSenderlUtil.sendMail(mail, subject, content, true, null, subject); + } + } else { + String content = subject + ",待消费消息数量:" + json; + maiSenderlUtil.sendMail(maiSenderlUtil.isEmail(focusMqCount.getSendMails()) ? focusMqCount.getSendMails() : "1129801211@qq.com", subject, content, true, null, subject); + } + } +// amqpService.SendMessage("mailNotice", json); + } + } } - } - if (queueMessageCountMap.size() > 0) { - //将queueMessageCountMap转成json字符串 - String json = JSON.toJSONString(queueMessageCountMap); - amqpService.SendMessage("mailNotice", json); + } catch (Exception e) { + e.printStackTrace(); + System.out.println("发送邮件失败"); + System.out.println(e.getMessage()); } } @@ -106,11 +144,11 @@ public class CheckRabbitMqScheduleTask { continue; } //3.发送邮件 - String content = "客户:" + focusDeviceOnline.getUserName() + ",设备离线数量:" + offlineDevices.size() + "个,"+"报警设备和预离线时间信息:"+json+"。请进入系统查看原因"; + String content = "客户:" + focusDeviceOnline.getUserName() + ",设备离线数量:" + offlineDevices.size() + "个," + "报警设备和预离线时间信息:" + json + "。请进入系统查看原因"; maiSenderlUtil.sendMail(mail, subject, content, true, null, subject); } } else { - String content = "客户:" + focusDeviceOnline.getUserName() + ",设备离线数量:" + offlineDevices.size() + "个,"+"报警设备和预离线时间信息:"+json+"。请进入系统查看原因"; + String content = "客户:" + focusDeviceOnline.getUserName() + ",设备离线数量:" + offlineDevices.size() + "个," + "报警设备和预离线时间信息:" + json + "。请进入系统查看原因"; maiSenderlUtil.sendMail(maiSenderlUtil.isEmail(focusDeviceOnline.getSendMails()) ? focusDeviceOnline.getSendMails() : "1129801211@qq.com", subject, content, true, null, subject); } } diff --git a/src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml b/src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml index 8210d4d..aa01503 100644 --- a/src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml +++ b/src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml @@ -25,6 +25,7 @@ + @@ -50,6 +51,7 @@ forward_time, last_source_data, last_forward_content, + top_forward_from, last_forward_time FROM topsail_transmit_log WHERE imei = #{imei} LIMIT 1 @@ -98,6 +100,14 @@ FROM topsail_product_data WHERE is_product_data = #{isProductData} + @@ -107,12 +117,12 @@ source_data, forward_content, forward_user_id, forward_user_name, forward_url, forward_port, forward_protocol, protocol, forward_result, forward_desc, forward_result_status, - forward_time, last_source_data, last_forward_content, last_forward_time) + forward_time, last_source_data, last_forward_content, last_forward_time, top_forward_from) VALUES (#{platform}, #{platformProtocol}, #{deviceId}, #{imei}, #{deviceType}, #{deviceTypeName}, #{imsi}, #{sourceData}, #{forwardContent}, #{forwardUserId}, #{forwardUserName}, #{forwardUrl}, #{forwardPort}, #{forwardProtocol}, #{protocol}, #{forwardResult}, #{forwardDesc}, #{forwardResultStatus}, - #{forwardTime}, #{lastSourceData}, #{lastForwardContent}, #{lastForwardTime}) + #{forwardTime}, #{lastSourceData}, #{lastForwardContent}, #{lastForwardTime}, #{topForwardFrom}) @@ -138,6 +148,7 @@ forward_time = #{forwardTime}, last_source_data = #{lastSourceData}, last_forward_content = #{lastForwardContent}, + top_forward_from = #{topForwardFrom}, last_forward_time = #{lastForwardTime} WHERE imei = #{imei}