This commit is contained in:
zhaohu
2026-06-12 16:14:58 +08:00
parent cf71897528
commit a7ba0f21e5
235 changed files with 33574 additions and 0 deletions
+24
View File
@@ -0,0 +1,24 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springblade</groupId>
<artifactId>blade-example</artifactId>
<version>${revision}</version>
</parent>
<artifactId>blade-mqtt-client</artifactId>
<name>${project.artifactId}</name>
<packaging>jar</packaging>
<dependencies>
<dependency>
<groupId>net.dreamlu</groupId>
<artifactId>mica-mqttx-client</artifactId>
<version>3.1.12</version>
</dependency>
</dependencies>
</project>
@@ -0,0 +1,47 @@
/**
* BladeX Commercial License Agreement
* Copyright (c) 2018-2099, https://bladex.cn. All rights reserved.
* <p>
* Use of this software is governed by the Commercial License Agreement
* obtained after purchasing a license from BladeX.
* <p>
* 1. This software is for development use only under a valid license
* from BladeX.
* <p>
* 2. Redistribution of this software's source code to any third party
* without a commercial license is strictly prohibited.
* <p>
* 3. Licensees may copyright their own code but cannot use segments
* from this software for such purposes. Copyright of this software
* remains with BladeX.
* <p>
* Using this software signifies agreement to this License, and the software
* must not be used for illegal purposes.
* <p>
* THIS SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY. The author is
* not liable for any claims arising from secondary or illegal development.
* <p>
* Author: Chill Zhuang (bladejava@qq.com)
*/
package org.springblade.mqtt.client;
import org.springblade.core.launch.BladeApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.scheduling.annotation.EnableScheduling;
import static org.springblade.common.constant.LauncherConstant.APPLICATION_MQTT_CLIENT_NAME;
/**
* MQTT Client 启动器
*
* @author Chill
*/
@EnableScheduling
@SpringBootApplication
public class MqttClientApplication {
public static void main(String[] args) {
BladeApplication.run(APPLICATION_MQTT_CLIENT_NAME, MqttClientApplication.class, args);
}
}
@@ -0,0 +1,98 @@
/**
* BladeX Commercial License Agreement
* Copyright (c) 2018-2099, https://bladex.cn. All rights reserved.
* <p>
* Use of this software is governed by the Commercial License Agreement
* obtained after purchasing a license from BladeX.
* <p>
* 1. This software is for development use only under a valid license
* from BladeX.
* <p>
* 2. Redistribution of this software's source code to any third party
* without a commercial license is strictly prohibited.
* <p>
* 3. Licensees may copyright their own code but cannot use segments
* from this software for such purposes. Copyright of this software
* remains with BladeX.
* <p>
* Using this software signifies agreement to this License, and the software
* must not be used for illegal purposes.
* <p>
* THIS SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY. The author is
* not liable for any claims arising from secondary or illegal development.
* <p>
* Author: Chill Zhuang (bladejava@qq.com)
*/
package org.springblade.mqtt.client.listener;
import lombok.extern.slf4j.Slf4j;
import org.springblade.mqtt.client.simulator.DeviceSimulator;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.stereotype.Component;
/**
* 设备应用监听器
*
* @author Chill
*/
@Slf4j
@Component
public class DeviceListener {
@Value("${mqtt.client.productKey:DEMO_PRODUCT}")
private String productKey;
@Value("${mqtt.client.deviceName:DEMO_DEVICE}")
private String deviceName;
@Value("${mqtt.client.deviceSecret:demo_secret_123}")
private String deviceSecret;
@Value("${mqtt.client.serverHost:localhost}")
private String serverHost;
@Value("${mqtt.client.serverPort:1883}")
private int serverPort;
@Value("${mqtt.client.enabled:true}")
private boolean enabled;
@EventListener(ApplicationReadyEvent.class)
public void handleApplicationReady() {
if (!enabled) {
log.info("MQTT客户端已禁用,跳过启动");
return;
}
log.info("应用启动完成,开始初始化MQTT客户端");
log.info("配置信息: productKey={}, deviceName={}, serverHost={}:{}",
productKey, deviceName, serverHost, serverPort);
try {
// 初始化设备模拟器
DeviceSimulator simulator = new DeviceSimulator(
productKey,
deviceName,
deviceSecret,
serverHost,
serverPort
);
// 启动模拟器
simulator.start();
log.info("设备模拟器启动成功");
// 添加关闭钩子
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
log.info("应用关闭,停止设备模拟器");
simulator.stop();
}));
} catch (Exception e) {
log.error("启动设备模拟器失败", e);
}
}
}
@@ -0,0 +1,282 @@
/**
* BladeX Commercial License Agreement
* Copyright (c) 2018-2099, https://bladex.cn. All rights reserved.
* <p>
* Use of this software is governed by the Commercial License Agreement
* obtained after purchasing a license from BladeX.
* <p>
* 1. This software is for development use only under a valid license
* from BladeX.
* <p>
* 2. Redistribution of this software's source code to any third party
* without a commercial license is strictly prohibited.
* <p>
* 3. Licensees may copyright their own code but cannot use segments
* from this software for such purposes. Copyright of this software
* remains with BladeX.
* <p>
* Using this software signifies agreement to this License, and the software
* must not be used for illegal purposes.
* <p>
* THIS SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY. The author is
* not liable for any claims arising from secondary or illegal development.
* <p>
* Author: Chill Zhuang (bladejava@qq.com)
*/
package org.springblade.mqtt.client.simulator;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import net.dreamlu.iot.mqtt.core.client.MqttClient;
import org.springblade.mqtt.client.support.DataReq;
import org.springblade.mqtt.client.support.NtpReq;
import org.tio.utils.buffer.ByteBufferUtil;
import org.tio.utils.json.JsonUtil;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
/**
* 设备模拟器
*
* @author Chill
*/
@Slf4j
@RequiredArgsConstructor
public class DeviceSimulator {
private final String productKey;
private final String deviceName;
private final String deviceSecret;
private final String serverHost;
private final int serverPort;
private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(4);
/**
* 启动模拟器
*/
public void start() {
log.info("启动设备模拟器: productKey={}, deviceName={}", productKey, deviceName);
// 初始化 MQTT 客户端
String clientId = generateClientId();
String username = generateUsername();
String password = generatePassword();
log.info("连接配置: clientId={}, username={}, host={}:{}", clientId, username, serverHost, serverPort);
try {
MqttClient client = MqttClient.create()
.ip(serverHost)
.port(serverPort)
.username(username)
.password(password)
.clientId(clientId)
.connectSync();
log.info("MQTT客户端连接成功");
// 订阅主题
subscribeTopics(client);
// 启动任务调度
scheduleTasks(client);
} catch (Exception e) {
log.error("MQTT客户端连接失败", e);
}
}
/**
* 生成客户端ID
*
* @return clientId
*/
private String generateClientId() {
return productKey + "_" + deviceName + "_" + System.currentTimeMillis();
}
/**
* 生成用户名
*
* @return username
*/
private String generateUsername() {
return deviceName + "&" + productKey;
}
/**
* 生成密码
*
* @return password
*/
private String generatePassword() {
// 简化的密码生成逻辑
return deviceSecret;
}
/**
* 订阅主题
*
* @param client MqttClient
*/
private void subscribeTopics(MqttClient client) {
// 订阅设备相关主题
String topicPrefix = "/blade/sys/" + productKey + "/" + deviceName;
// 订阅属性设置
client.subQos0(topicPrefix + "/thing/service/property/set", (context, topic, message, payload) -> {
log.info("收到属性设置: topic={}, payload={}", topic, ByteBufferUtil.toString(payload));
});
// 订阅服务调用
client.subQos0(topicPrefix + "/thing/service/+", (context, topic, message, payload) -> {
log.info("收到服务调用: topic={}, payload={}", topic, ByteBufferUtil.toString(payload));
});
// 订阅NTP响应
client.subQos0("/blade/ext/ntp/" + productKey + "/" + deviceName + "/response", (context, topic, message, payload) -> {
log.info("收到NTP响应: topic={}, payload={}", topic, ByteBufferUtil.toString(payload));
});
log.info("已订阅设备主题: {}", topicPrefix);
}
/**
* 任务调度
*
* @param client MqttClient
*/
private void scheduleTasks(MqttClient client) {
// 定期上报属性数据(每30秒)
scheduler.scheduleAtFixedRate(() -> {
publishPropertyData(client);
}, 5, 30, TimeUnit.SECONDS);
// 定期上报事件数据(每60秒)
scheduler.scheduleAtFixedRate(() -> {
publishEventData(client);
}, 10, 60, TimeUnit.SECONDS);
// 定期发送心跳数据(每120秒)
scheduler.scheduleAtFixedRate(() -> {
publishHeartbeat(client);
}, 15, 120, TimeUnit.SECONDS);
// 定期请求NTP时间(每300秒)
scheduler.scheduleAtFixedRate(() -> {
publishNtpRequest(client);
}, 20, 300, TimeUnit.SECONDS);
log.info("已启动定时任务调度");
}
/**
* 发布属性数据
*
* @param client MqttClient
*/
private void publishPropertyData(MqttClient client) {
try {
DataReq<Map<String, Object>> req = new DataReq<>();
req.setId(UUID.randomUUID().toString().replace("-", ""));
req.setVersion("1.0");
Map<String, Object> params = new HashMap<>();
params.put("temperature", 20 + Math.random() * 10); // 温度:20-30度
params.put("humidity", 40 + Math.random() * 30); // 湿度:40-70%
params.put("lightSwitch", Math.random() > 0.5 ? 1 : 0); // 灯开关
req.setParams(params);
String topic = "/blade/sys/" + productKey + "/" + deviceName + "/thing/event/property/post";
client.publish(topic, JsonUtil.toJsonBytes(req));
log.info("已发布属性数据: {}", JsonUtil.toJsonString(params));
} catch (Exception e) {
log.error("发布属性数据失败", e);
}
}
/**
* 发布事件数据
*
* @param client MqttClient
*/
private void publishEventData(MqttClient client) {
try {
DataReq<Map<String, Object>> req = new DataReq<>();
req.setId(UUID.randomUUID().toString().replace("-", ""));
req.setVersion("1.0");
Map<String, Object> params = new HashMap<>();
Map<String, Object> output = new HashMap<>();
output.put("eventLevel", "INFO");
output.put("eventMessage", "设备运行正常");
output.put("timestamp", System.currentTimeMillis());
params.put("output", JsonUtil.toJsonString(output));
params.put("eventName", "状态事件");
params.put("eventType", "info");
req.setParams(params);
String topic = "/blade/sys/" + productKey + "/" + deviceName + "/thing/event/StatusEvent/post";
client.publish(topic, JsonUtil.toJsonBytes(req));
log.info("已发布事件数据: {}", JsonUtil.toJsonString(params));
} catch (Exception e) {
log.error("发布事件数据失败", e);
}
}
/**
* 发布心跳数据
*
* @param client MqttClient
*/
private void publishHeartbeat(MqttClient client) {
try {
Map<String, Object> heartbeat = new HashMap<>();
heartbeat.put("deviceId", deviceName);
heartbeat.put("productKey", productKey);
heartbeat.put("timestamp", System.currentTimeMillis());
heartbeat.put("status", "online");
String topic = "/blade/sys/" + productKey + "/" + deviceName + "/heartbeat";
client.publish(topic, JsonUtil.toJsonBytes(heartbeat));
log.info("已发布心跳数据: {}", JsonUtil.toJsonString(heartbeat));
} catch (Exception e) {
log.error("发布心跳数据失败", e);
}
}
/**
* 发布NTP请求
*
* @param client MqttClient
*/
private void publishNtpRequest(MqttClient client) {
try {
NtpReq req = new NtpReq();
req.setDeviceSendTime(String.valueOf(System.currentTimeMillis()));
String topic = "/blade/ext/ntp/" + productKey + "/" + deviceName + "/request";
client.publish(topic, JsonUtil.toJsonBytes(req));
log.info("已发布NTP请求: {}", JsonUtil.toJsonString(req));
} catch (Exception e) {
log.error("发布NTP请求失败", e);
}
}
/**
* 停止模拟器
*/
public void stop() {
scheduler.shutdown();
log.info("设备模拟器已停止");
}
}
@@ -0,0 +1,67 @@
/**
* BladeX Commercial License Agreement
* Copyright (c) 2018-2099, https://bladex.cn. All rights reserved.
* <p>
* Use of this software is governed by the Commercial License Agreement
* obtained after purchasing a license from BladeX.
* <p>
* 1. This software is for development use only under a valid license
* from BladeX.
* <p>
* 2. Redistribution of this software's source code to any third party
* without a commercial license is strictly prohibited.
* <p>
* 3. Licensees may copyright their own code but cannot use segments
* from this software for such purposes. Copyright of this software
* remains with BladeX.
* <p>
* Using this software signifies agreement to this License, and the software
* must not be used for illegal purposes.
* <p>
* THIS SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY. The author is
* not liable for any claims arising from secondary or illegal development.
* <p>
* Author: Chill Zhuang (bladejava@qq.com)
*/
package org.springblade.mqtt.client.support;
import com.fasterxml.jackson.annotation.JsonInclude;
import lombok.Data;
import java.io.Serializable;
/**
* 设备通用数据格式
*
* @author Chill
*/
@Data
@JsonInclude(JsonInclude.Include.NON_NULL)
public class DataReq<T> implements Serializable {
/**
* 消息ID号。uuid,去掉短横线,32位,全局唯一,用于ack或系统消息追踪
*/
private String id;
/**
* 协议版本号,目前协议版本号唯一取值为1.0
*/
private String version;
/**
* 扩展功能的参数,其下包含各功能字段。平台可扩展,或可自行扩展,自行扩展的参数需在自定义解析模块自行解析
*/
private SysBean sys;
/**
* 请求方法。
*/
private String method;
/**
* 请求参数
*/
private T params;
}
@@ -0,0 +1,45 @@
/**
* BladeX Commercial License Agreement
* Copyright (c) 2018-2099, https://bladex.cn. All rights reserved.
* <p>
* Use of this software is governed by the Commercial License Agreement
* obtained after purchasing a license from BladeX.
* <p>
* 1. This software is for development use only under a valid license
* from BladeX.
* <p>
* 2. Redistribution of this software's source code to any third party
* without a commercial license is strictly prohibited.
* <p>
* 3. Licensees may copyright their own code but cannot use segments
* from this software for such purposes. Copyright of this software
* remains with BladeX.
* <p>
* Using this software signifies agreement to this License, and the software
* must not be used for illegal purposes.
* <p>
* THIS SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY. The author is
* not liable for any claims arising from secondary or illegal development.
* <p>
* Author: Chill Zhuang (bladejava@qq.com)
*/
package org.springblade.mqtt.client.support;
import lombok.Data;
import java.io.Serializable;
/**
* NTP 请求
*
* @author Chill
*/
@Data
public class NtpReq implements Serializable {
/**
* 设备端发送时间
*/
private String deviceSendTime;
}
@@ -0,0 +1,62 @@
/**
* BladeX Commercial License Agreement
* Copyright (c) 2018-2099, https://bladex.cn. All rights reserved.
* <p>
* Use of this software is governed by the Commercial License Agreement
* obtained after purchasing a license from BladeX.
* <p>
* 1. This software is for development use only under a valid license
* from BladeX.
* <p>
* 2. Redistribution of this software's source code to any third party
* without a commercial license is strictly prohibited.
* <p>
* 3. Licensees may copyright their own code but cannot use segments
* from this software for such purposes. Copyright of this software
* remains with BladeX.
* <p>
* Using this software signifies agreement to this License, and the software
* must not be used for illegal purposes.
* <p>
* THIS SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY. The author is
* not liable for any claims arising from secondary or illegal development.
* <p>
* Author: Chill Zhuang (bladejava@qq.com)
*/
package org.springblade.mqtt.client.support;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.io.Serializable;
/**
* 扩展功能
*
* @author Chill
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class SysBean implements Serializable {
/**
* 不响应 ack
*/
public static final int ACK_NO = 0;
/**
* 响应 ack
*/
public static final int ACK_NEED = 1;
public SysBean(boolean ackNeed) {
this(ackNeed ? ACK_NEED : ACK_NO);
}
/**
* 扩展功能字段,表示是否返回响应数据。0:不返回响应数据1:返回响应数据
*/
private int ack;
}
@@ -0,0 +1,29 @@
# 服务器配置
server:
port: 8291
undertow:
# 以下的配置会影响buffer,这些buffer会用于服务器连接的IO操作,有点类似netty的池化内存管理
buffer-size: 1024
# 是否分配的直接内存
direct-buffers: true
# 线程配置
threads:
# 设置IO线程数, 它主要执行非阻塞的任务,它们会负责多个连接, 默认设置每个CPU核心一个线程
io: 16
# 阻塞任务线程池, 当执行类似servlet请求阻塞操作, undertow会从这个线程池中取得线程,它的值设置取决于系统的负载
worker: 400
servlet:
# 编码配置
encoding:
charset: UTF-8
force: true
# MQTT客户端配置
mqtt:
client:
enabled: true # 是否启用MQTT客户端
productKey: "BLADE_DEMO_PRODUCT" # 产品标识
deviceName: "BLADE_DEMO_DEVICE_001" # 设备名称
deviceSecret: "blade_demo_secret_123456" # 设备密钥
serverHost: "localhost" # MQTT服务器地址
serverPort: 11883 # MQTT服务器端口