Browse Source

1.优化rabbitMQ存储消息方法2.信息设备转发日志存储功能

master
bgy 3 months ago
parent
commit
d1c8d8ee49
13 changed files with 414 additions and 83 deletions
  1. +51
    -0
      src/main/java/com/topsail/scheduletask/config/RabbitMQConfig.java
  2. +40
    -0
      src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java
  3. +82
    -0
      src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java
  4. +73
    -0
      src/main/java/com/topsail/scheduletask/receiver/AmqpListener.java
  5. +41
    -66
      src/main/java/com/topsail/scheduletask/result/Result.java
  6. +2
    -2
      src/main/java/com/topsail/scheduletask/service/MqMonitorService.java
  7. +1
    -1
      src/main/java/com/topsail/scheduletask/task/HeartBeatTask.java
  8. +1
    -1
      src/main/java/com/topsail/scheduletask/task/ServiceDownMonitor.java
  9. +0
    -6
      src/main/resources/application-prod.properties
  10. +0
    -6
      src/main/resources/application-test.properties
  11. +70
    -0
      src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml
  12. +26
    -0
      src/main/resources/db/migration/V20260512__add_fields_to_topsail_transmit_log.sql
  13. +27
    -1
      src/main/resources/schema.sql

+ 51
- 0
src/main/java/com/topsail/scheduletask/config/RabbitMQConfig.java View File

@ -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);
}
}

+ 40
- 0
src/main/java/com/topsail/scheduletask/mapper/TopsailTransmitLogDao.java View File

@ -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 <pre>
* Author Version Date Changes
* Administrator 1.0 2021年11月24日 Created
*
* </pre>
* @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);
}

+ 82
- 0
src/main/java/com/topsail/scheduletask/pojo/TopsailTransmitLog.java View File

@ -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;
}

+ 73
- 0
src/main/java/com/topsail/scheduletask/receiver/AmqpListener.java View File

@ -4,6 +4,8 @@ import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.rabbitmq.client.Channel; import com.rabbitmq.client.Channel;
import com.topsail.scheduletask.mapper.InformLogDao; 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.service.AmqpService;
import com.topsail.scheduletask.util.MaiSenderlUtil; import com.topsail.scheduletask.util.MaiSenderlUtil;
import org.slf4j.Logger; import org.slf4j.Logger;
@ -29,6 +31,8 @@ public class AmqpListener {
AmqpService amqpService; AmqpService amqpService;
@Autowired @Autowired
InformLogDao informLogDao; InformLogDao informLogDao;
@Autowired
TopsailTransmitLogDao topsailTransmitLogDao;
@RabbitListener(queues = "mailNotice") @RabbitListener(queues = "mailNotice")
public void sendMailMessage(@Payload String message, @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag, Channel channel) throws Exception { 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);
}
}
}
} }

+ 41
- 66
src/main/java/com/topsail/scheduletask/result/Result.java View File

@ -1,89 +1,64 @@
package com.topsail.scheduletask.result; package com.topsail.scheduletask.result;
import com.github.pagehelper.PageInfo;
import lombok.Data; import lombok.Data;
import org.springframework.validation.BindingResult; import org.springframework.validation.BindingResult;
import org.springframework.validation.FieldError; 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; import java.util.stream.Collectors;
@Data @Data
public class Result<T> { public class Result<T> {
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<T> Result<T> success(T data){
return new Result<T>(data);
}
/**
* 失败时候的调用
* */
public static <T> Result<T> error(CodeMsg codeMsg){
return new Result<T>(codeMsg);
}
private Result(CodeMsg codeMsg) {
if (codeMsg != null) {
this.code = codeMsg.getCode();
this.msg = codeMsg.getMsg();
}
}
public static <T> Result<T> error(){
return new Result<T>(CodeMsg.FAILED);
}
/**
* 成功时候的调用
*/
public static <T> Result<T> success(T data) {
return new Result<T>(data);
}
/**
* 失败时候的调用
*/
public static <T> Result<T> error(CodeMsg codeMsg) {
return new Result<T>(codeMsg);
}
public static <T> Result<T> error() {
return new Result<T>(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 <T>
* @return
*/
public static<T> Map<String,Object> returnPageMap(List<T> t){
PageInfo<T> pageInfo=new PageInfo(t);
int count=(int)pageInfo.getTotal();
Map<String,Object> map=new HashMap<>();
map.put("list", new ArrayList(pageInfo.getList()));
map.put("count",count);
return map;
}
} }

+ 2
- 2
src/main/java/com/topsail/scheduletask/service/MqMonitorService.java View File

@ -34,8 +34,8 @@ public class MqMonitorService {
private static final Map<String, String> MONITOR_QUEUES = new HashMap<>(); private static final Map<String, String> MONITOR_QUEUES = new HashMap<>();
static { 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<String, Long> lastAlertTime = new ConcurrentHashMap<>(); private Map<String, Long> lastAlertTime = new ConcurrentHashMap<>();


+ 1
- 1
src/main/java/com/topsail/scheduletask/task/HeartBeatTask.java View File

@ -18,7 +18,7 @@ public class HeartBeatTask {
@Autowired @Autowired
private HeartBeatMapper heartBeatMapper; private HeartBeatMapper heartBeatMapper;
@Scheduled(fixedRate = 10000)
// @Scheduled(fixedRate = 10000)
public void heartbeat() { public void heartbeat() {
try { try {
String serviceName = ApplicationContextUtil.getApplicationName(); String serviceName = ApplicationContextUtil.getApplicationName();


+ 1
- 1
src/main/java/com/topsail/scheduletask/task/ServiceDownMonitor.java View File

@ -36,7 +36,7 @@ public class ServiceDownMonitor {
private Map<String, Long> lastAlertTime = new ConcurrentHashMap<>(); private Map<String, Long> lastAlertTime = new ConcurrentHashMap<>();
private static final long ALERT_INTERVAL = 60000; private static final long ALERT_INTERVAL = 60000;
@Scheduled(fixedRate = 60000)
// @Scheduled(fixedRate = 60000)
public void checkServiceDown() { public void checkServiceDown() {
try { try {
List<ServiceHeartbeat> deadList = heartBeatMapper.selectTimeoutService(timeout); List<ServiceHeartbeat> deadList = heartBeatMapper.selectTimeoutService(timeout);


+ 0
- 6
src/main/resources/application-prod.properties View File

@ -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

+ 0
- 6
src/main/resources/application-test.properties View File

@ -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

+ 70
- 0
src/main/resources/com/topsail/scheduletask/mapper/TopsailTransmitLogMapper.xml View File

@ -0,0 +1,70 @@
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.topsail.scheduletask.mapper.TopsailTransmitLogDao">
<resultMap id="BaseResultMap" type="com.topsail.scheduletask.pojo.TopsailTransmitLog">
<result column="platform" property="platform"/>
<result column="device_id" property="deviceId"/>
<result column="imei" property="imei"/>
<result column="device_type" property="deviceType"/>
<result column="device_type_name" property="deviceTypeName"/>
<result column="imsi" property="imsi"/>
<result column="source_data" property="sourceData"/>
<result column="forward_content" property="forwardContent"/>
<result column="forward_user_id" property="forwardUserId"/>
<result column="forward_user_name" property="forwardUserName"/>
<result column="forward_url" property="forwardUrl"/>
<result column="forward_port" property="forwardPort"/>
<result column="forward_protocol" property="forwardProtocol"/>
<result column="forward_result" property="forwardResult"/>
<result column="forward_time" property="forwardTime"/>
<result column="last_forward_time" property="lastForwardTime"/>
</resultMap>
<!-- 根据 IMEI 查询记录 -->
<select id="getTopsailTransmitLogByImei" resultMap="BaseResultMap">
SELECT 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
FROM topsail_transmit_log
WHERE imei = #{imei}
LIMIT 1
</select>
<!-- 新增记录 -->
<insert id="insertTopsailTransmitLog" parameterType="com.topsail.scheduletask.pojo.TopsailTransmitLog">
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}
)
</insert>
<!-- 根据 IMEI 更新记录 -->
<update id="updateTopsailTransmitLogByImei" parameterType="com.topsail.scheduletask.pojo.TopsailTransmitLog">
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}
</update>
</mapper>

+ 26
- 0
src/main/resources/db/migration/V20260512__add_fields_to_topsail_transmit_log.sql View File

@ -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`;

+ 27
- 1
src/main/resources/schema.sql View File

@ -36,4 +36,30 @@ CREATE TABLE IF NOT EXISTS `alert_record` (
PRIMARY KEY (`id`), PRIMARY KEY (`id`),
KEY `idx_alert_type` (`alert_type`), KEY `idx_alert_type` (`alert_type`),
KEY `idx_status` (`status`) KEY `idx_status` (`status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='告警记录表';
) 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='设备数据转发日志表';

Loading…
Cancel
Save