diff --git a/pom.xml b/pom.xml
index 8eab379..b0db7cc 100644
--- a/pom.xml
+++ b/pom.xml
@@ -40,6 +40,12 @@
1.18.30
+
+
+ org.eclipse.paho
+ org.eclipse.paho.client.mqttv3
+ 1.2.5
+
com.aliyun
diff --git a/src/main/java/com/project/operation/domain/entity/OperationLogEntity.java b/src/main/java/com/project/operation/domain/entity/OperationLogEntity.java
index 9488a7d..899abd4 100644
--- a/src/main/java/com/project/operation/domain/entity/OperationLogEntity.java
+++ b/src/main/java/com/project/operation/domain/entity/OperationLogEntity.java
@@ -53,4 +53,10 @@ public class OperationLogEntity extends BaseEntity {
@Column(name = "client_ip", columnDefinition="varchar(100) comment '客户端IP'")
@TableField("client_ip")
private String clientIp;
+
+
+ @Column(name = "is_reported")
+ @Comment("是否已上报MQTT, 0-否, 1-是")
+ @TableField("is_reported")
+ private Integer isReported = 0;
}
diff --git a/src/main/java/com/project/operation/domain/service/MqttLogReporterService.java b/src/main/java/com/project/operation/domain/service/MqttLogReporterService.java
new file mode 100644
index 0000000..3527371
--- /dev/null
+++ b/src/main/java/com/project/operation/domain/service/MqttLogReporterService.java
@@ -0,0 +1,211 @@
+package com.project.operation.domain.service;
+
+import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
+import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.project.operation.domain.entity.OperationLogEntity;
+import com.project.operation.mapper.OperationLogMapper;
+import jakarta.annotation.PostConstruct;
+import jakarta.annotation.PreDestroy;
+import org.eclipse.paho.client.mqttv3.MqttClient;
+import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
+import org.eclipse.paho.client.mqttv3.MqttMessage;
+import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Service;
+
+
+import java.net.NetworkInterface;
+import java.nio.charset.StandardCharsets;
+import java.text.SimpleDateFormat;
+import java.util.Date;
+import java.util.Enumeration;
+import java.util.HashMap;
+import java.util.Map;
+
+@Service
+public class MqttLogReporterService {
+
+ private static final String BROKER_URL = "tcp://oam.itc-pa.cn:1883";
+ private static final String USERNAME = "ychjasp";
+ private static final String PASSWORD = "821006";
+
+ // 依照您的要求,固定配置参数
+ private static final int PROJECT_ID = 2282; // 项目 ID (int)
+ private static final String MODEL = "ES0450"; // 产品型号 (String)
+ private static final String COMPANY = "BL"; // 公司名称 (固定)
+ private static final String DEVICE_NAME = "8.129.84.155"; // 设备名称 (固定)
+ private static final String PRODUCT_NAME = "智考云"; // 产品线名称 (固定)
+
+ @Autowired
+ private OperationLogMapper operationLogMapper;
+
+ private MqttClient mqttClient;
+ private String macAddress;
+ private String pubTopic;
+
+ // 用于安全生成嵌套 JSON 的工具
+ private final ObjectMapper objectMapper = new ObjectMapper();
+
+ @PostConstruct
+ public void init() {
+ this.macAddress = getLocalMacAddress();
+ this.pubTopic = "device/" + this.macAddress + "/report/log";
+ connect();
+ }
+
+ private void connect() {
+ try {
+ mqttClient = new MqttClient(BROKER_URL, macAddress, new MemoryPersistence());
+ MqttConnectOptions options = new MqttConnectOptions();
+ options.setUserName(USERNAME);
+ options.setPassword(PASSWORD.toCharArray());
+ options.setCleanSession(true);
+ options.setAutomaticReconnect(true);
+ options.setKeepAliveInterval(20);
+ options.setConnectionTimeout(30);
+
+ mqttClient.connect(options);
+ System.out.println("MQTT 连接成功。ClientID: " + macAddress);
+ } catch (Exception e) {
+ System.err.println("MQTT 连接失败: " + e.getMessage());
+ }
+ }
+
+ /**
+ * 定时扫描任务:每 11 秒执行一次。
+ * 只捞取一条未上报(is_reported = 0)的日志,保障上报间隔不小于10秒
+ */
+ @Scheduled(fixedDelay = 11000)
+ public void scanAndReport() {
+ if (mqttClient == null || !mqttClient.isConnected()) {
+ return;
+ }
+
+ try {
+ // 1. 获取最早的一条未上报数据
+ LambdaQueryWrapper queryWrapper = new LambdaQueryWrapper<>();
+ queryWrapper.eq(OperationLogEntity::getIsReported, 0)
+ .orderByAsc(OperationLogEntity::getId)
+ .last("LIMIT 1");
+
+ OperationLogEntity logEntity = operationLogMapper.selectOne(queryWrapper);
+
+ // 2. 如果存在未上报数据,将其上报并更新状态
+ if (logEntity != null) {
+ String payload = buildJsonPayload(logEntity);
+ MqttMessage message = new MqttMessage(payload.getBytes(StandardCharsets.UTF_8));
+ message.setQos(1);
+
+ // 3. 发布到 MQTT
+ mqttClient.publish(pubTopic, message);
+
+ // 4. 更新状态为 1
+ logEntity.setIsReported(1);
+ operationLogMapper.updateById(logEntity);
+
+ System.out.println("日志上报成功,ID: " + logEntity.getId() + " 报文内容: " + payload);
+ }
+ } catch (Exception e) {
+ System.err.println("日志定时上报异常: " + e.getMessage());
+ }
+ }
+
+ /**
+ * 严格对照您提供的接口字段及嵌套格式生成 JSON
+ */
+ private String buildJsonPayload(OperationLogEntity logEntity) throws Exception {
+
+ // 1. 组装内层的 data 对象
+ Map dataMap = new HashMap<>();
+
+ // 字段: timestamp (设备端事件发生的时间戳,秒级)
+ long timestamp = logEntity.getCreateTime() != null ? logEntity.getCreateTime().getTime() / 1000 : System.currentTimeMillis() / 1000;
+ dataMap.put("timestamp", (int) timestamp);
+
+ // 字段: log_type (1--异常 2--业务 3--运行 4--配置)
+ // 映射规则:如果是失败(result=1)或有异常信息,归为 1 (异常),否则归为 2 (业务)
+ int logType = 2;
+ dataMap.put("log_type", logType);
+
+ // 字段: log_level (1--轻微 2--一般 3--严重 4--致命,正常操作发0)
+ int logLevel = 1;
+ dataMap.put("log_level", logLevel);
+
+ // 字段: log_desc (日志描述,不超过1024字节)
+ String logDesc = String.format("[%s-%s] %s | 方法: %s | 耗时: %dms",
+ logEntity.getModule(),
+ logEntity.getAction(),
+ logEntity.getDescription(),
+ logEntity.getMethod(),
+ logEntity.getCostTime()
+ );
+ if (logDesc.length() > 300) { // 预防超长,做个截断保护
+ logDesc = logDesc.substring(0, 300) + "...";
+ }
+ dataMap.put("log_desc", logDesc);
+
+ // 字段: serial_number (唯一序列号,不超过100字符)
+ // 映射规则:直接使用数据库该条日志的主键 ID 作为唯一标识
+ dataMap.put("serial_number", String.valueOf(logEntity.getId()));
+
+ // 字段: product_name (产品线名称,不超过50字符)
+ dataMap.put("product_name", PRODUCT_NAME);
+
+ // 字段: model (型号名称)
+ dataMap.put("model", MODEL);
+
+ // 字段: log_source (日志维度 1-11)
+ // 映射规则:程序报错归为 7,正常操作归为 4(其他用户操作日志)
+ int logSource = (logType == 1) ? 7 : 1;
+ dataMap.put("log_source", logSource);
+
+ // 字段: project_id (项目id,必选)
+ dataMap.put("project_id", PROJECT_ID);
+
+ // 2. 组装外层对象
+ Map rootMap = new HashMap<>();
+ rootMap.put("company", COMPANY);
+ rootMap.put("device_name", DEVICE_NAME);
+ rootMap.put("data", dataMap);
+
+ // 3. 序列化为标准 JSON
+ return objectMapper.writeValueAsString(rootMap);
+ }
+
+ private String getLocalMacAddress() {
+ try {
+ Enumeration interfaces = NetworkInterface.getNetworkInterfaces();
+ while (interfaces.hasMoreElements()) {
+ NetworkInterface ni = interfaces.nextElement();
+ if (ni.isLoopback() || ni.isVirtual() || !ni.isUp()) {
+ continue;
+ }
+ byte[] mac = ni.getHardwareAddress();
+ if (mac != null && mac.length == 6) {
+ StringBuilder sb = new StringBuilder();
+ for (int i = 0; i < mac.length; i++) {
+ sb.append(String.format("%02X%s", mac[i], (i < mac.length - 1) ? "-" : ""));
+ }
+ return sb.toString();
+ }
+ }
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+ return "00-11-22-33-44-55";
+ }
+
+ @PreDestroy
+ public void cleanup() {
+ try {
+ if (mqttClient != null && mqttClient.isConnected()) {
+ mqttClient.disconnect();
+ mqttClient.close();
+ }
+ } catch (Exception e) {
+ // silent
+ }
+ }
+}