From d1c8d8ee49223a016d1e8c79a7989b8dba6ef8b4 Mon Sep 17 00:00:00 2001 From: bgy Date: Tue, 12 May 2026 15:50:40 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=BC=98=E5=8C=96rabbitMQ=E5=AD=98=E5=82=A8?= =?UTF-8?q?=E6=B6=88=E6=81=AF=E6=96=B9=E6=B3=952.=E4=BF=A1=E6=81=AF?= =?UTF-8?q?=E8=AE=BE=E5=A4=87=E8=BD=AC=E5=8F=91=E6=97=A5=E5=BF=97=E5=AD=98?= =?UTF-8?q?=E5=82=A8=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../scheduletask/config/RabbitMQConfig.java | 51 ++++++++++ .../scheduletask/mapper/TopsailTransmitLogDao.java | 40 ++++++++ .../scheduletask/pojo/TopsailTransmitLog.java | 82 ++++++++++++++++ .../scheduletask/receiver/AmqpListener.java | 73 ++++++++++++++ .../com/topsail/scheduletask/result/Result.java | 107 ++++++++------------- .../scheduletask/service/MqMonitorService.java | 4 +- .../topsail/scheduletask/task/HeartBeatTask.java | 2 +- .../scheduletask/task/ServiceDownMonitor.java | 2 +- src/main/resources/application-prod.properties | 6 -- src/main/resources/application-test.properties | 6 -- .../mapper/TopsailTransmitLogMapper.xml | 70 ++++++++++++++ ...0260512__add_fields_to_topsail_transmit_log.sql | 26 +++++ src/main/resources/schema.sql | 28 +++++- 13 files changed, 414 insertions(+), 83 deletions(-) create mode 100644 src/main/java/com/topsail/scheduletask/config/RabbitMQConfig.java create mode 100644 src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java create mode 100644 src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java delete mode 100644 src/main/resources/application-prod.properties delete mode 100644 src/main/resources/application-test.properties create mode 100644 src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml create mode 100644 src/main/resources/db/migration/V20260512__add_fields_to_topsail_transmit_log.sql diff --git a/src/main/java/com/topsail/scheduletask/config/RabbitMQConfig.java b/src/main/java/com/topsail/scheduletask/config/RabbitMQConfig.java new file mode 100644 index 0000000..8896c15 --- /dev/null +++ b/src/main/java/com/topsail/scheduletask/config/RabbitMQConfig.java @@ -0,0 +1,51 @@ +package com.topsail.scheduletask.config; + +import org.springframework.amqp.core.AmqpAdmin; +import org.springframework.amqp.core.AmqpTemplate; +import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; +import org.springframework.amqp.rabbit.connection.ConnectionFactory; +import org.springframework.amqp.rabbit.core.RabbitAdmin; +import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +@Configuration +public class RabbitMQConfig { + + @Value("${spring.rabbitmq.host:localhost}") + private String host; + + @Value("${spring.rabbitmq.port:5672}") + private int port; + + @Value("${spring.rabbitmq.username:guest}") + private String username; + + @Value("${spring.rabbitmq.password:guest}") + private String password; + + @Value("${spring.rabbitmq.virtualHost:/}") + private String virtualHost; + + @Bean + public ConnectionFactory connectionFactory() { + CachingConnectionFactory connectionFactory = new CachingConnectionFactory(); + connectionFactory.setHost(host); + connectionFactory.setPort(port); + connectionFactory.setUsername(username); + connectionFactory.setPassword(password); + connectionFactory.setVirtualHost(virtualHost); + return connectionFactory; + } + + @Bean + public AmqpAdmin amqpAdmin(ConnectionFactory connectionFactory) { + return new RabbitAdmin(connectionFactory); + } + + @Bean + public AmqpTemplate amqpTemplate(ConnectionFactory connectionFactory) { + return new RabbitTemplate(connectionFactory); + } +} \ No newline at end of file diff --git a/src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java b/src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java new file mode 100644 index 0000000..01f9b5b --- /dev/null +++ b/src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java @@ -0,0 +1,40 @@ +package com.topsail.scheduletask.mapper; + +import com.topsail.scheduletask.pojo.TopsailTransmitLog; +import org.apache.ibatis.annotations.Param; +import org.springframework.stereotype.Repository; + +/** + * 信息推送日志表Dao + * + * @version
+ * Author	Version		Date		Changes
+ * Administrator 	1.0  		2021年11月24日 	Created
+ *
+ * 
+ * @since 1. + */ +@Repository +public interface TopsailTransmitLogDao { + + /** + * 根据 IMEI 查询记录 + * @param imei 设备号 + * @return TopsailTransmitLog 对象 + */ + TopsailTransmitLog getTopsailTransmitLogByImei(@Param("imei") String imei); + + /** + * 新增记录 + * @param topsailTransmitLog 待新增的对象 + * @return 影响行数 + */ + int insertTopsailTransmitLog(TopsailTransmitLog topsailTransmitLog); + + /** + * 根据 IMEI 更新记录 + * @param topsailTransmitLog 待更新的对象 + * @return 影响行数 + */ + int updateTopsailTransmitLogByImei(TopsailTransmitLog topsailTransmitLog); +} diff --git a/src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java b/src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java new file mode 100644 index 0000000..5bc5da8 --- /dev/null +++ b/src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java @@ -0,0 +1,82 @@ +package com.topsail.scheduletask.pojo; + +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.NoArgsConstructor; + + +/** + * 通用明文数据转发对象 + */ +@Data +@AllArgsConstructor +@NoArgsConstructor +@Builder +public class TopsailTransmitLog { + //1.平台信息 + /** + * 平台 + */ + private String platform; + //2.设备基础信息 + /** + * 设备ID + */ + private String deviceId; + /** + * 设备号 + */ + private String imei; + /** + * 设备类型 + */ + private Integer deviceType; + private String deviceTypeName; + + /** + * imsi + */ + private String imsi; + /** + * 源数据 + */ + private String sourceData; + /** + * 转发内容 + */ + private String forwardContent; + /*** + * 转发用户id + */ + private Integer forwardUserId; + /** + * 转发用户名称 + */ + private String forwardUserName; + /** + * 转发url + */ + private String forwardUrl; + /** + * 转发端口 + */ + private Integer forwardPort; + /** + * 转发协议类型 + */ + private String forwardProtocol; + /** + * 转发结果 + */ + private String forwardResult; + /** + * 转发时间 + */ + private String forwardTime; + /** + * 上一次转发时间 + */ + private String lastForwardTime; + +} diff --git a/src/main/java/com/topsail/scheduletask/receiver/AmqpListener.java b/src/main/java/com/topsail/scheduletask/receiver/AmqpListener.java index 62fbdbf..2d7ed64 100644 --- a/src/main/java/com/topsail/scheduletask/receiver/AmqpListener.java +++ b/src/main/java/com/topsail/scheduletask/receiver/AmqpListener.java @@ -4,6 +4,8 @@ import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSONObject; import com.rabbitmq.client.Channel; import com.topsail.scheduletask.mapper.InformLogDao; +import com.topsail.scheduletask.mapper.TopsailTransmitLogDao; +import com.topsail.scheduletask.pojo.TopsailTransmitLog; import com.topsail.scheduletask.service.AmqpService; import com.topsail.scheduletask.util.MaiSenderlUtil; import org.slf4j.Logger; @@ -29,6 +31,8 @@ public class AmqpListener { AmqpService amqpService; @Autowired InformLogDao informLogDao; + @Autowired + TopsailTransmitLogDao topsailTransmitLogDao; @RabbitListener(queues = "mailNotice") public void sendMailMessage(@Payload String message, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag, Channel channel) throws Exception { @@ -51,5 +55,74 @@ public class AmqpListener { } } + @RabbitListener(queues = "topsailtransmitlog") + public void saveTopsailTransmitLog(@Payload String message, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag, Channel channel) { + try { + // 1. 参数校验 + if (message == null || message.trim().isEmpty()) { + LOG.warn("接收到空消息,跳过处理"); + channel.basicAck(deliveryTag, false); + return; + } + + // 2. 解析消息 + JSONObject jsonObject = JSONObject.parseObject(message); + if (jsonObject == null) { + LOG.error("消息解析失败: {}", message); + channel.basicAck(deliveryTag, false); + return; + } + + TopsailTransmitLog topsailTransmitLog = JSON.toJavaObject(jsonObject, TopsailTransmitLog.class); + if (topsailTransmitLog == null) { + LOG.error("转换为 TopsailTransmitLog 对象失败: {}", message); + channel.basicAck(deliveryTag, false); + return; + } + + String imei = topsailTransmitLog.getImei(); + if (imei == null || imei.trim().isEmpty()) { + LOG.warn("IMEI 为空,跳过处理: {}", message); + channel.basicAck(deliveryTag, false); + return; + } + + // 3. 查询是否存在记录 + TopsailTransmitLog oldTopsailTransmitLog = topsailTransmitLogDao.getTopsailTransmitLogByImei(imei); + + // 4. 根据情况执行新增或更新 + if (oldTopsailTransmitLog == null) { + // 新增记录 + int insertResult = topsailTransmitLogDao.insertTopsailTransmitLog(topsailTransmitLog); + if (insertResult > 0) { + LOG.debug("成功新增 TopsailTransmitLog 记录, IMEI: {}", imei); + } else { + LOG.error("新增 TopsailTransmitLog 记录失败, IMEI: {}", imei); + } + } else { + // 更新记录 + topsailTransmitLog.setLastForwardTime(oldTopsailTransmitLog.getForwardTime()); + int updateResult = topsailTransmitLogDao.updateTopsailTransmitLogByImei(topsailTransmitLog); + if (updateResult > 0) { + LOG.debug("成功更新 TopsailTransmitLog 记录, IMEI: {}", imei); + } else { + LOG.error("更新 TopsailTransmitLog 记录失败, IMEI: {}", imei); + } + } + + // 5. 手动确认消息 + channel.basicAck(deliveryTag, false); + + } catch (Exception e) { + LOG.error("处理 topsailtransmitlog 消息时发生异常: {}", message, e); + try { + // 发生异常时拒绝消息,并重新入队 + channel.basicNack(deliveryTag, false, true); + } catch (Exception ex) { + LOG.error("拒绝消息时发生异常", ex); + } + } + } + } diff --git a/src/main/java/com/topsail/scheduletask/result/Result.java b/src/main/java/com/topsail/scheduletask/result/Result.java index e826a5d..035e983 100644 --- a/src/main/java/com/topsail/scheduletask/result/Result.java +++ b/src/main/java/com/topsail/scheduletask/result/Result.java @@ -1,89 +1,64 @@ package com.topsail.scheduletask.result; -import com.github.pagehelper.PageInfo; import lombok.Data; import org.springframework.validation.BindingResult; import org.springframework.validation.FieldError; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; import java.util.stream.Collectors; @Data public class Result { - - private int code; - private String msg="success"; - private T data; - private Result(T data) { - this.data = data; - } + private int code; + private String msg = "success"; + private T data; - private Result(int code, String msg) { - this.code = code; - this.msg = msg; - } + private Result(T data) { + this.data = data; + } - private Result(CodeMsg codeMsg) { - if(codeMsg != null) { - this.code = codeMsg.getCode(); - this.msg = codeMsg.getMsg(); - } - } + private Result(int code, String msg) { + this.code = code; + this.msg = msg; + } - /** - * 成功时候的调用 - * */ - public static Result success(T data){ - return new Result(data); - } - - /** - * 失败时候的调用 - * */ - public static Result error(CodeMsg codeMsg){ - return new Result(codeMsg); - } + private Result(CodeMsg codeMsg) { + if (codeMsg != null) { + this.code = codeMsg.getCode(); + this.msg = codeMsg.getMsg(); + } + } - public static Result error(){ - return new Result(CodeMsg.FAILED); - } - + /** + * 成功时候的调用 + */ + public static Result success(T data) { + return new Result(data); + } - + /** + * 失败时候的调用 + */ + public static Result error(CodeMsg codeMsg) { + return new Result(codeMsg); + } + public static Result error() { + return new Result(CodeMsg.FAILED); + } - /** - * BindingResult统一处理 - */ - public static Result resolveBindResult(BindingResult bindingResult){ - StringBuilder stringBuilder = new StringBuilder(); - for (String s : bindingResult.getFieldErrors().stream().map(FieldError::getDefaultMessage).collect(Collectors.toList())) { - stringBuilder.append(",").append(s); - } - return Result.error(new CodeMsg(502,stringBuilder.toString().substring(1))); - } + /** + * BindingResult统一处理 + */ + public static Result resolveBindResult(BindingResult bindingResult) { + StringBuilder stringBuilder = new StringBuilder(); + for (String s : bindingResult.getFieldErrors().stream().map(FieldError::getDefaultMessage).collect(Collectors.toList())) { + stringBuilder.append(",").append(s); + } + return Result.error(new CodeMsg(502, stringBuilder.toString().substring(1))); + } - /** - *分页之后的统一返回对象 - * @param t - * @param - * @return - */ - public static Map returnPageMap(List t){ - PageInfo pageInfo=new PageInfo(t); - int count=(int)pageInfo.getTotal(); - Map map=new HashMap<>(); - map.put("list", new ArrayList(pageInfo.getList())); - map.put("count",count); - return map; - } - - } diff --git a/src/main/java/com/topsail/scheduletask/service/MqMonitorService.java b/src/main/java/com/topsail/scheduletask/service/MqMonitorService.java index bb9b02c..b7ed04c 100644 --- a/src/main/java/com/topsail/scheduletask/service/MqMonitorService.java +++ b/src/main/java/com/topsail/scheduletask/service/MqMonitorService.java @@ -34,8 +34,8 @@ public class MqMonitorService { private static final Map MONITOR_QUEUES = new HashMap<>(); static { - MONITOR_QUEUES.put("order_queue", "order"); - MONITOR_QUEUES.put("pay_queue", "pay"); +// MONITOR_QUEUES.put("order_queue", "order"); +// MONITOR_QUEUES.put("pay_queue", "pay"); } private Map lastAlertTime = new ConcurrentHashMap<>(); diff --git a/src/main/java/com/topsail/scheduletask/task/HeartBeatTask.java b/src/main/java/com/topsail/scheduletask/task/HeartBeatTask.java index 312c505..472dfd2 100644 --- a/src/main/java/com/topsail/scheduletask/task/HeartBeatTask.java +++ b/src/main/java/com/topsail/scheduletask/task/HeartBeatTask.java @@ -18,7 +18,7 @@ public class HeartBeatTask { @Autowired private HeartBeatMapper heartBeatMapper; - @Scheduled(fixedRate = 10000) +// @Scheduled(fixedRate = 10000) public void heartbeat() { try { String serviceName = ApplicationContextUtil.getApplicationName(); diff --git a/src/main/java/com/topsail/scheduletask/task/ServiceDownMonitor.java b/src/main/java/com/topsail/scheduletask/task/ServiceDownMonitor.java index 7ef8b82..5732593 100644 --- a/src/main/java/com/topsail/scheduletask/task/ServiceDownMonitor.java +++ b/src/main/java/com/topsail/scheduletask/task/ServiceDownMonitor.java @@ -36,7 +36,7 @@ public class ServiceDownMonitor { private Map lastAlertTime = new ConcurrentHashMap<>(); private static final long ALERT_INTERVAL = 60000; - @Scheduled(fixedRate = 60000) +// @Scheduled(fixedRate = 60000) public void checkServiceDown() { try { List deadList = heartBeatMapper.selectTimeoutService(timeout); diff --git a/src/main/resources/application-prod.properties b/src/main/resources/application-prod.properties deleted file mode 100644 index 1bd8324..0000000 --- a/src/main/resources/application-prod.properties +++ /dev/null @@ -1,6 +0,0 @@ -spring.rabbitmq.host=182.92.218.150 -spring.rabbitmq.username=topsail -spring.rabbitmq.password=topsail -#spring.datasource.url=jdbc:mysql://localhost:3306/iot?useSSL=false&useUnicode=true&characterEncoding=utf-8 -#spring.datasource.username=topsail -#spring.datasource.password=topsail2020 \ No newline at end of file diff --git a/src/main/resources/application-test.properties b/src/main/resources/application-test.properties deleted file mode 100644 index 7cf8979..0000000 --- a/src/main/resources/application-test.properties +++ /dev/null @@ -1,6 +0,0 @@ -spring.rabbitmq.host=47.97.117.253 -#spring.rabbitmq.username=topsail -#spring.rabbitmq.password=top -#spring.datasource.url=jdbc:mysql://localhost:3306/iot?useSSL=false&useUnicode=true&characterEncoding=utf-8 -#spring.datasource.username=topsail -#spring.datasource.password=topsail2020 \ No newline at end of file diff --git a/src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml b/src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml new file mode 100644 index 0000000..89f3cd2 --- /dev/null +++ b/src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml @@ -0,0 +1,70 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + INSERT INTO topsail_transmit_log ( + platform, device_id, imei, device_type, device_type_name, imsi, + source_data, forward_content, forward_user_id, forward_user_name, + forward_url, forward_port, forward_protocol, forward_result, + forward_time, last_forward_time + ) VALUES ( + #{platform}, #{deviceId}, #{imei}, #{deviceType}, #{deviceTypeName}, #{imsi}, + #{sourceData}, #{forwardContent}, #{forwardUserId}, #{forwardUserName}, + #{forwardUrl}, #{forwardPort}, #{forwardProtocol}, #{forwardResult}, + #{forwardTime}, #{lastForwardTime} + ) + + + + + UPDATE topsail_transmit_log + SET platform = #{platform}, + device_id = #{deviceId}, + device_type = #{deviceType}, + device_type_name = #{deviceTypeName}, + imsi = #{imsi}, + source_data = #{sourceData}, + forward_content = #{forwardContent}, + forward_user_id = #{forwardUserId}, + forward_user_name = #{forwardUserName}, + forward_url = #{forwardUrl}, + forward_port = #{forwardPort}, + forward_protocol = #{forwardProtocol}, + forward_result = #{forwardResult}, + forward_time = #{forwardTime}, + last_forward_time = #{lastForwardTime} + WHERE imei = #{imei} + + + diff --git a/src/main/resources/db/migration/V20260512__add_fields_to_topsail_transmit_log.sql b/src/main/resources/db/migration/V20260512__add_fields_to_topsail_transmit_log.sql new file mode 100644 index 0000000..68217e6 --- /dev/null +++ b/src/main/resources/db/migration/V20260512__add_fields_to_topsail_transmit_log.sql @@ -0,0 +1,26 @@ +-- 为 topsail_transmit_log 表添加新字段 +-- 执行时间: 2026-05-12 + +-- 添加设备类型名称字段 +ALTER TABLE `topsail_transmit_log` +ADD COLUMN `device_type_name` varchar(64) DEFAULT NULL COMMENT '设备类型名称' AFTER `device_type`; + +-- 添加转发用户ID字段 +ALTER TABLE `topsail_transmit_log` +ADD COLUMN `forward_user_id` int(11) DEFAULT NULL COMMENT '转发用户ID' AFTER `forward_content`; + +-- 添加转发用户名称字段 +ALTER TABLE `topsail_transmit_log` +ADD COLUMN `forward_user_name` varchar(64) DEFAULT NULL COMMENT '转发用户名称' AFTER `forward_user_id`; + +-- 添加转发URL字段 +ALTER TABLE `topsail_transmit_log` +ADD COLUMN `forward_url` varchar(256) DEFAULT NULL COMMENT '转发URL' AFTER `forward_user_name`; + +-- 添加转发端口字段 +ALTER TABLE `topsail_transmit_log` +ADD COLUMN `forward_port` int(11) DEFAULT NULL COMMENT '转发端口' AFTER `forward_url`; + +-- 添加转发协议类型字段 +ALTER TABLE `topsail_transmit_log` +ADD COLUMN `forward_protocol` varchar(32) DEFAULT NULL COMMENT '转发协议类型' AFTER `forward_port`; diff --git a/src/main/resources/schema.sql b/src/main/resources/schema.sql index 32b58a7..6c77de3 100644 --- a/src/main/resources/schema.sql +++ b/src/main/resources/schema.sql @@ -36,4 +36,30 @@ CREATE TABLE IF NOT EXISTS `alert_record` ( PRIMARY KEY (`id`), KEY `idx_alert_type` (`alert_type`), KEY `idx_status` (`status`) -) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='告警记录表'; \ No newline at end of file +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='告警记录表'; + +CREATE TABLE IF NOT EXISTS `topsail_transmit_log` ( + `id` bigint(20) NOT NULL AUTO_INCREMENT, + `platform` varchar(64) DEFAULT NULL COMMENT '平台', + `device_id` varchar(64) DEFAULT NULL COMMENT '设备ID', + `imei` varchar(64) NOT NULL COMMENT '设备号', + `device_type` int(11) DEFAULT NULL COMMENT '设备类型', + `device_type_name` varchar(64) DEFAULT NULL COMMENT '设备类型名称', + `imsi` varchar(64) DEFAULT NULL COMMENT 'IMSI', + `source_data` text COMMENT '源数据', + `forward_content` text COMMENT '转发内容', + `forward_user_id` int(11) DEFAULT NULL COMMENT '转发用户ID', + `forward_user_name` varchar(64) DEFAULT NULL COMMENT '转发用户名称', + `forward_url` varchar(256) DEFAULT NULL COMMENT '转发URL', + `forward_port` int(11) DEFAULT NULL COMMENT '转发端口', + `forward_protocol` varchar(32) DEFAULT NULL COMMENT '转发协议类型', + `forward_result` varchar(512) DEFAULT NULL COMMENT '转发结果', + `forward_time` varchar(32) DEFAULT NULL COMMENT '转发时间', + `last_forward_time` varchar(32) DEFAULT NULL COMMENT '上一次转发时间', + `create_time` datetime DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间', + `update_time` datetime DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间', + PRIMARY KEY (`id`), + UNIQUE KEY `uk_imei` (`imei`), + KEY `idx_device_id` (`device_id`), + KEY `idx_forward_time` (`forward_time`) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='设备数据转发日志表'; \ No newline at end of file