物联网设备接入与通信
设备接入认证
物联网平台接入认证的核心目标是确保只有合法设备能够接入平台,防止仿冒和伪造设备。常见的认证方式包括一机一密、一型一密和 X.509 证书认证。
一机一密
一机一密指每台设备拥有唯一的设备密钥(DeviceSecret),设备在连接平台时使用该密钥对通信参数进行签名,平台端校验签名合法性后放行连接。
认证流程
设备 IoT 平台
| |
|-- 1. 携带 ProductKey + DeviceName + ClientId + Timestamp -->|
| |
| 2. 平台根据 ProductKey + DeviceName |
| 查询对应的 DeviceSecret |
| |
|<-- 3. 返回 Challenge / 非对称应答 (或不返回,由设备主动签名)-|
| |
|-- 4. 使用 DeviceSecret 对参数签名: |
| sign = HMACSHA1(DeviceSecret, params) |
| |
| 5. 平台用相同 DeviceSecret 计算签名 |
| 并与设备上传 signature 比对 |
| |
|<-- 6. 认证通过,建立 MQTT/TCP 连接 ---------------------------|HMAC 签名示例
// 设备端签名算法
public class DeviceAuthUtil {
public static String sign(String deviceSecret, String params) throws Exception {
Mac mac = Mac.getInstance("HmacSHA1");
SecretKeySpec key = new SecretKeySpec(
deviceSecret.getBytes(StandardCharsets.UTF_8), "HmacSHA1");
mac.init(key);
byte[] digest = mac.doFinal(params.getBytes(StandardCharsets.UTF_8));
return Hex.encodeHexString(digest).toUpperCase();
}
// 待签名字符串: "clientId123deviceNameMyDeviceproductKeya1b2c3timestamp1710000000"
public static String buildSignContent(String clientId, String deviceName,
String productKey, long timestamp) {
return "clientId" + clientId
+ "deviceName" + deviceName
+ "productKey" + productKey
+ "timestamp" + timestamp;
}
}建表模型
-- 设备认证信息表
CREATE TABLE device_auth (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
product_key VARCHAR(32) NOT NULL COMMENT '产品 key',
device_name VARCHAR(64) NOT NULL COMMENT '设备名称',
device_secret VARCHAR(64) NOT NULL COMMENT '设备密钥',
status TINYINT NOT NULL DEFAULT 0 COMMENT '0-未激活 1-已激活 2-已禁用',
active_time DATETIME COMMENT '首次激活时间',
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
gmt_modified DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_product_device (product_key, device_name)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='设备认证信息';一型一密
一型一密指同一型号(ProductKey)下的设备共享一个 ProductSecret,设备首次上线时通过 ProductSecret 完成自动注册并获得唯一的 DeviceSecret,后续使用 DeviceSecret 做一机一密认证。
自动注册流程
设备 IoT 平台
| |
|-- 1. 携带 ProductKey + DeviceName + ProductSecret 签名 -->|
| |
| 2. 校验 ProductSecret 合法性 |
| 3. 自动生成 DeviceSecret |
| 4. 写入 device_auth 表 |
| |
|<-- 5. 返回 DeviceSecret (仅首次注册时返回) -------------|
| |
|-- 6. 后续连接使用一机一密流程 (使用已分配的 DeviceSecret) |一型一密配置
# 产品密钥配置
iot:
product:
enabled: true
auto-register: true # 开启自动注册
max-devices-per-product: 10000 # 单产品最大设备数
secret-ttl-seconds: 86400 # ProductSecret 有效性缓存时间@Service
public class DeviceAutoRegisterService {
@Autowired
private DeviceAuthMapper deviceAuthMapper;
public DeviceRegisterResult autoRegister(String productKey, String deviceName,
String clientId, long timestamp, String signature) {
// 1. 从缓存/数据库获取 ProductSecret
String productSecret = getProductSecret(productKey);
// 2. 校验签名
String expectedSign = sign(productSecret,
"clientId" + clientId + "deviceName" + deviceName +
"productKey" + productKey + "timestamp" + timestamp);
if (!expectedSign.equals(signature)) {
throw new AuthException("product secret sign mismatch");
}
// 3. 检查设备是否已注册
DeviceAuth device = deviceAuthMapper.selectByProductKeyAndDeviceName(productKey, deviceName);
if (device != null) {
// 已注册,返回已有 DeviceSecret
return DeviceRegisterResult.existed(device.getDeviceSecret());
}
// 4. 自动注册,生成新的 DeviceSecret
String newDeviceSecret = generateDeviceSecret();
deviceAuthMapper.insert(new DeviceAuth(productKey, deviceName, newDeviceSecret));
return DeviceRegisterResult.newly(newDeviceSecret);
}
}X.509 证书认证
X.509 证书认证基于 PKI(公钥基础设施)体系,由 CA 为每个设备签发数字证书,通过 TLS 双向认证(mTLS)完成设备身份核验。适用于安全等级要求较高的场景,如工业物联网、车联网。
认证架构
CA 根证书 (自签名)
|
+---------+---------+
| |
中间 CA 证书 中间 CA 证书
| |
+---+---+ +----+----+
| | | |
设备A 设备B 设备C 设备DTLS 双向认证握手
证书管理
-- 设备证书表
CREATE TABLE device_certificate (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
product_key VARCHAR(32) NOT NULL,
device_name VARCHAR(64) NOT NULL,
cert_serial VARCHAR(128) NOT NULL COMMENT '证书序列号',
cert_cn VARCHAR(128) NOT NULL COMMENT '证书通用名称',
issuer_dn VARCHAR(256) NOT NULL COMMENT '颁发者 DN',
cert_pem TEXT NOT NULL COMMENT '证书 PEM 内容',
private_key_pem TEXT COMMENT '私钥 PEM (仅首次返回)',
cert_status TINYINT NOT NULL DEFAULT 0 COMMENT '0-有效 1-吊销 2-过期',
not_before DATETIME NOT NULL COMMENT '证书生效时间',
not_after DATETIME NOT NULL COMMENT '证书过期时间',
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_cert_serial (cert_serial),
UNIQUE KEY uk_product_device (product_key, device_name)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='设备证书表';
-- 证书吊销列表
CREATE TABLE cert_revocation_list (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
cert_serial VARCHAR(128) NOT NULL,
revoked_at DATETIME NOT NULL COMMENT '吊销时间',
reason VARCHAR(64) COMMENT '吊销原因',
crl_publish_at DATETIME NOT NULL COMMENT 'CRL 发布时间'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='证书吊销列表';Spring Boot mTLS 配置
server:
ssl:
enabled: true
client-auth: need # 双向认证
key-store: classpath:iot-platform.p12
key-store-password: ${KEYSTORE_PASS}
key-store-type: PKCS12
trust-store: classpath:trust-store.p12
trust-store-password: ${TRUSTSTORE_PASS}
trust-store-type: PKCS12
port: 8443证书校验示例
@Component
public class X509AuthFilter implements Filter {
@Override
public void doFilter(ServletRequest req, ServletResponse res, FilterChain chain)
throws IOException, ServletException {
HttpServletRequest request = (HttpServletRequest) req;
X509Certificate[] certs = (X509Certificate[]) request.getAttribute(
"jakarta.servlet.request.X509Certificate");
if (certs == null || certs.length == 0) {
throw new AuthenticationException("missing client certificate");
}
X509Certificate clientCert = certs[0];
try {
// 验证证书未过期
clientCert.checkValidity();
// 验证证书是否已吊销 (CRL/OCSP)
checkRevocation(clientCert);
// 提取设备身份信息
String deviceName = extractCN(clientCert.getSubjectX500Principal());
String certSerial = clientCert.getSerialNumber().toString();
// 校验设备名与证书的绑定关系
validateDeviceCertBinding(deviceName, certSerial);
} catch (CertPathValidatorException | CRLException e) {
throw new AuthenticationException("certificate validation failed", e);
}
chain.doFilter(req, res);
}
}设备注册流程
设备注册支持三种模式:预注册、自动注册、批量导入。
预注册
管理员在平台侧预先录入设备信息(ProductKey、DeviceName、DeviceSecret),设备出厂时写入凭证,首次上线直接认证。
@PostMapping("/api/v1/devices/pre-register")
public Result<Void> preRegister(@Valid @RequestBody PreRegisterReq req) {
// 校验产品是否存在
Product product = productService.getByProductKey(req.getProductKey());
Assert.notNull(product, "product not found");
// 检查设备名是否已存在
DeviceAuth exist = deviceAuthMapper.selectByProductKeyAndDeviceName(
req.getProductKey(), req.getDeviceName());
Assert.isNull(exist, "device already registered");
// 写入设备认证信息
DeviceAuth device = new DeviceAuth();
device.setProductKey(req.getProductKey());
device.setDeviceName(req.getDeviceName());
device.setDeviceSecret(generateDeviceSecret());
device.setStatus(DeviceStatus.INACTIVE);
deviceAuthMapper.insert(device);
return Result.success(device);
}自动注册
设备使用一型一密方式首次上线时由平台自动完成注册。需要控制自动注册的速率和总设备上限,避免被恶意刷量。
iot:
auto-register:
enabled: true
rate-limit: 100 # 每分钟每产品最大自动注册数
max-devices: 100000 # 单产品最大设备数
ip-whitelist: # 可选注册 IP 白名单
- 10.0.0.0/8
- 172.16.0.0/12批量导入
通过 CSV/Excel 文件批量导入设备,适用于量产场景。
public class DeviceBatchImportService {
public BatchImportResult importFromCsv(MultipartFile file, String productKey) {
List<DeviceAuth> devices = new ArrayList<>();
try (BufferedReader reader = new BufferedReader(
new InputStreamReader(file.getInputStream(), StandardCharsets.UTF_8))) {
String line;
while ((line = reader.readLine()) != null) {
String[] fields = line.split(",");
if (fields.length < 1) continue;
DeviceAuth device = new DeviceAuth();
device.setProductKey(productKey);
device.setDeviceName(fields[0].trim());
device.setDeviceSecret(generateDeviceSecret());
device.setStatus(DeviceStatus.INACTIVE);
devices.add(device);
}
}
// 批量写入,使用 ignore 避免重复
deviceAuthMapper.batchInsertIgnore(devices);
return BatchImportResult.of(devices.size());
}
}-- 批量写入 (MySQL)
INSERT IGNORE INTO device_auth (product_key, device_name, device_secret, status)
VALUES
('pk1', 'dev001', 'sec001', 0),
('pk1', 'dev002', 'sec002', 0),
('pk1', 'dev003', 'sec003', 0);三种注册方式对比
| 模式 | 适用场景 | 安全性 | 运维复杂度 | 出厂流程 |
|---|---|---|---|---|
| 预注册 | 高安全要求、固定设备 | 最高 | 高 | 需预烧录 DeviceSecret |
| 自动注册 | 消费级设备、快速量产 | 中等 | 低 | 只需烧录 ProductSecret |
| 批量导入 | 企业级设备批量接入 | 高 | 中 | 烧录批量导出的密钥 |
设备影子
设备影子(Device Shadow)是一个 JSON 文档,用于缓存设备的最新状态。设备离线时平台存储其最新状态,设备上线后自动同步。应用层可以读取影子获取设备状态,也可以更新影子中的 desired 状态指挥设备。
影子文档结构
{
"state": {
"reported": {
"temperature": 25.6,
"humidity": 68.2,
"switch": "on",
"firmware_ver": "2.1.0"
},
"desired": {
"temperature": 26.0,
"switch": "off"
}
},
"metadata": {
"reported": {
"temperature": { "timestamp": 1710000001 },
"humidity": { "timestamp": 1710000001 },
"switch": { "timestamp": 1710000002 },
"firmware_ver":{ "timestamp": 1710000000 }
},
"desired": {
"temperature": { "timestamp": 1710000100 },
"switch": { "timestamp": 1710000100 }
}
},
"version": 42,
"timestamp": 1710000100
}| 字段 | 说明 |
|---|---|
state.reported | 设备主动上报的状态(设备 -> 平台) |
state.desired | 应用层期望设备达到的状态(平台 -> 设备) |
metadata | 每个属性的更新时间戳 |
version | 影子文档版本号,递增且用于乐观锁冲突检测 |
timestamp | 影子文档最后更新时间 |
影子文档存储
CREATE TABLE device_shadow (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
product_key VARCHAR(32) NOT NULL,
device_name VARCHAR(64) NOT NULL,
shadow_doc JSON NOT NULL COMMENT '影子文档 JSON',
version BIGINT NOT NULL DEFAULT 0 COMMENT '影子版本号',
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
gmt_modified DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_product_device (product_key, device_name)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='设备影子';影子状态同步机制
状态更新流程
设备 IoT 平台 应用
| | |
|-- MQTT: report state ------>| |
| topic: /shadow/update | |
| payload: {"temperature": | |
| 25.6} | |
| | |
| 1. 校验版本号 (乐观锁) |
| 2. 合并 reported 字段 |
| 3. version++ |
| 4. 比对 desired,若有差异则下发 desired |
| |
|<-- MQTT: desired delta ----| |
| topic: /shadow/update | |
| payload: {"switch":"off"}| |
| |
|-- MQTT: update delta ack ->| |
| | |
| |-- HTTP GET /shadow ------->|
| |<-- shadow JSON ------------|设备上报状态
// 设备端上报 MQTT 消息
public class ShadowReporter {
private static final String SHADOW_UPDATE_TOPIC = "/shadow/update";
public void reportTemperature(double temperature) {
JsonObject reported = new JsonObject();
reported.addProperty("temperature", temperature);
reported.addProperty("timestamp", System.currentTimeMillis() / 1000);
JsonObject state = new JsonObject();
state.add("reported", reported);
JsonObject payload = new JsonObject();
payload.add("state", state);
payload.addProperty("version", shadowVersion); // 当前影子版本号
mqttClient.publish(SHADOW_UPDATE_TOPIC, payload.toString().getBytes(), 1);
}
}平台端影子更新
@Service
public class ShadowUpdateService {
@Autowired
private DeviceShadowMapper shadowMapper;
@Autowired
private RedisTemplate<String, String> redisTemplate;
private static final String SHADOW_LOCK_PREFIX = "shadow:lock:";
@Transactional
public boolean updateShadow(String productKey, String deviceName,
JsonObject reported, Long expectedVersion) {
String lockKey = SHADOW_LOCK_PREFIX + productKey + ":" + deviceName;
RLock lock = redissonClient.getLock(lockKey);
try {
// 分布式锁防止并发更新
if (!lock.tryLock(3, 10, TimeUnit.SECONDS)) {
throw new ShadowUpdateException("acquire shadow lock timeout");
}
DeviceShadow shadow = shadowMapper.selectByProductKeyAndDeviceName(
productKey, deviceName);
// 乐观锁版本校验
if (expectedVersion != null && !expectedVersion.equals(shadow.getVersion())) {
throw new ShadowVersionConflictException(
"version conflict, expected: " + expectedVersion
+ ", actual: " + shadow.getVersion());
}
// 合并 reported 字段
JsonObject currentDoc = JsonParser.parseString(shadow.getShadowDoc()).getAsJsonObject();
JsonObject currentReported = currentDoc.getAsJsonObject("state")
.getAsJsonObject("reported");
mergeJson(currentReported, reported);
updateMetadata(currentDoc, "reported", reported);
// 版本递增
long newVersion = shadow.getVersion() + 1;
currentDoc.addProperty("version", newVersion);
currentDoc.addProperty("timestamp", System.currentTimeMillis() / 1000);
// 更新数据库
shadowMapper.updateShadowDoc(productKey, deviceName,
currentDoc.toString(), newVersion, shadow.getVersion());
// 检查 desired 是否需要下发
checkAndDeliverDesired(productKey, deviceName, currentDoc);
return true;
} finally {
if (lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
}
}desired 下发
private void checkAndDeliverDesired(String productKey, String deviceName, JsonObject shadowDoc) {
JsonObject state = shadowDoc.getAsJsonObject("state");
JsonObject desired = state.getAsJsonObject("desired");
JsonObject reported = state.getAsJsonObject("reported");
if (desired == null || desired.size() == 0) {
return;
}
// 计算 delta: desired - reported
JsonObject delta = new JsonObject();
for (Map.Entry<String, JsonElement> entry : desired.entrySet()) {
String key = entry.getKey();
if (reported == null || !reported.has(key)
|| !reported.get(key).equals(entry.getValue())) {
delta.add(key, entry.getValue());
}
}
if (delta.size() > 0) {
// 通过 MQTT 下发 delta 到设备
String topic = "/shadow/delta/" + productKey + "/" + deviceName;
JsonObject deltaMsg = new JsonObject();
deltaMsg.add("delta", delta);
deltaMsg.addProperty("version", shadowDoc.get("version").getAsLong());
mqttGateway.sendToTopic(topic, deltaMsg.toString());
}
}影子缓存优化
影子文档频繁读写,直接操作数据库会有性能瓶颈。通过引入 Redis 缓存减少数据库压力,同时可以缓存设备最近的影子状态减少设备端请求。
缓存架构
设备/应用 Redis (缓存层) MySQL (持久层)
| | |
|-- 读影子 --------------->| |
|<-- 返回 (LRU 缓存) -----| |
| | |
|-- 写影子 --------------->| |
| |-- 异步写 behind (延迟双写) ---->|
| | (Redis 作为主存储) |
| | |
| | 定期全量持久化 / binlog 同步 |
| |<---------------------------------|缓存实现
@Component
public class ShadowCacheManager {
private static final String SHADOW_CACHE_KEY = "shadow:%s:%s";
private static final Duration CACHE_TTL = Duration.ofHours(2);
@Autowired
private StringRedisTemplate redisTemplate;
@Autowired
private DeviceShadowMapper shadowMapper;
public JsonObject getShadow(String productKey, String deviceName) {
String cacheKey = String.format(SHADOW_CACHE_KEY, productKey, deviceName);
// 1. 查缓存
String cached = redisTemplate.opsForValue().get(cacheKey);
if (cached != null) {
return JsonParser.parseString(cached).getAsJsonObject();
}
// 2. 缓存 miss,查数据库
DeviceShadow shadow = shadowMapper.selectByProductKeyAndDeviceName(
productKey, deviceName);
if (shadow == null) {
return createEmptyShadow(productKey, deviceName);
}
// 3. 回填缓存
JsonObject doc = JsonParser.parseString(shadow.getShadowDoc()).getAsJsonObject();
redisTemplate.opsForValue().set(cacheKey, doc.toString(), CACHE_TTL);
return doc;
}
@Transactional
public void updateShadowAndCache(String productKey, String deviceName,
JsonObject reported, Long version) {
String cacheKey = String.format(SHADOW_CACHE_KEY, productKey, deviceName);
// 1. 更新数据库 (事务内)
shadowMapper.updateReported(productKey, deviceName, reported.toString(), version);
// 2. 更新缓存 (缓存为主,DB 异步持久化的场景可使用 Redis 作为主写入)
JsonObject doc = getShadow(productKey, deviceName);
mergeIntoReported(doc, reported);
doc.addProperty("version", doc.get("version").getAsLong() + 1);
redisTemplate.opsForValue().set(cacheKey, doc.toString(), CACHE_TTL);
}
}断连恢复优化
设备离线期间影子持续缓存最新状态,设备上线后一次性推送离线期间的 desired delta,而非逐条推送。
// 设备上线时批量同步离线期间的 desired 变更
public void onDeviceConnected(String productKey, String deviceName) {
JsonObject shadow = shadowCacheManager.getShadow(productKey, deviceName);
JsonObject state = shadow.getAsJsonObject("state");
if (state == null) return;
JsonObject desired = state.getAsJsonObject("desired");
JsonObject reported = state.getAsJsonObject("reported");
if (desired == null || desired.size() == 0) return;
// 计算完整 delta 一次性下发
JsonObject delta = new JsonObject();
for (Map.Entry<String, JsonElement> entry : desired.entrySet()) {
String key = entry.getKey();
if (reported == null || !reported.has(key)
|| !reported.get(key).equals(entry.getValue())) {
delta.add(key, entry.getValue());
}
}
if (delta.size() > 0) {
JsonObject deltaMsg = new JsonObject();
deltaMsg.add("delta", delta);
deltaMsg.addProperty("version", shadow.get("version").getAsLong());
mqttGateway.sendToTopic(
"/shadow/delta/" + productKey + "/" + deviceName,
deltaMsg.toString());
}
}数据上报
设备端采集数据后上报到平台,平台完成解析、校验、存储和转发。数据上报涉及消息格式、解析脚本、频率控制和时序数据存储等环节。
消息格式规范
平台支持多种消息格式,设备可根据自身资源情况选择。
JSON 格式
{
"id": "msg_1710000001",
"method": "thing.event.property.post",
"params": {
"temperature": 25.6,
"humidity": 68.2,
"pressure": 1013.2
},
"timestamp": 1710000001000
}二进制格式
字段定义:
Bytes 0-1: 消息类型 (0x01 = 属性上报)
Bytes 2-3: 消息体长度 (N)
Bytes 4-N+3: 消息体 (TLV 编码)
Bytes N+4: 校验和 (CRC8)
TLV 编码示例:
Tag=0x01 (temperature) Length=4 Value=0x41CC0000 (25.6 float)
Tag=0x02 (humidity) Length=4 Value=0x42886666 (68.2 float)Protobuf 格式
syntax = "proto3";
package iot.device;
message PropertyReport {
string msg_id = 1;
int64 timestamp = 2;
float temperature = 3;
float humidity = 4;
float pressure = 5;
repeated TagValue tags = 10;
}
message TagValue {
string key = 1;
string value = 2;
}TLV 编解码示例
public class TlvCodec {
public static byte[] encode(List<TlvEntry> entries) {
ByteBuffer buf = ByteBuffer.allocate(1024);
buf.order(ByteOrder.BIG_ENDIAN);
for (TlvEntry entry : entries) {
buf.putShort(entry.getTag());
buf.putShort((short) entry.getValue().length);
buf.put(entry.getValue());
}
buf.flip();
byte[] result = new byte[buf.remaining()];
buf.get(result);
return result;
}
public static List<TlvEntry> decode(byte[] data) {
List<TlvEntry> entries = new ArrayList<>();
ByteBuffer buf = ByteBuffer.wrap(data);
buf.order(ByteOrder.BIG_ENDIAN);
while (buf.remaining() >= 4) {
short tag = buf.getShort();
short len = buf.getShort();
byte[] val = new byte[len];
buf.get(val);
entries.add(new TlvEntry(tag, val));
}
return entries;
}
}数据解析脚本
对于二进制、TLV 等非标准格式,平台提供脚本引擎进行解析,将原始字节转换为平台统一的 JSON 格式。
脚本引擎架构
设备上报二进制数据
|
v
协议适配层 (识别 Protocol Type)
|
v
脚本引擎 (Groovy/JS/Nashorn)
- 从脚本仓库加载对应脚本
- 沙箱执行 (禁止网络/文件操作)
- 输入: byte[] / Map<String,Object>
- 输出: JSON (统一格式)
|
v
平台统一处理 (校验/存储/转发)脚本管理
CREATE TABLE parse_script (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
product_key VARCHAR(32) NOT NULL,
script_name VARCHAR(128) NOT NULL COMMENT '脚本名称',
script_lang VARCHAR(16) NOT NULL DEFAULT 'groovy' COMMENT 'groovy/js',
script_content MEDIUMTEXT NOT NULL COMMENT '脚本内容',
script_version INT NOT NULL DEFAULT 1,
status TINYINT NOT NULL DEFAULT 0 COMMENT '0-草稿 1-已发布 2-已下线',
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
gmt_modified DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_product_script (product_key, script_name, script_version)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='解析脚本';Groovy 脚本示例
// 二进制温湿度传感器解析脚本
import java.nio.ByteBuffer
import java.nio.ByteOrder
def parse(byte[] rawData) {
ByteBuffer buf = ByteBuffer.wrap(rawData)
buf.order(ByteOrder.LITTLE_ENDIAN)
def result = [:]
result.msgId = "raw_" + System.currentTimeMillis()
result.method = "thing.event.property.post"
result.timestamp = System.currentTimeMillis()
def params = [:]
// 读取温度 (2 bytes, 0.1 精度)
short tempRaw = buf.getShort()
params.temperature = tempRaw / 10.0
// 读取湿度 (2 bytes, 0.1 精度)
short humRaw = buf.getShort()
params.humidity = humRaw / 10.0
// 读取开关状态 (1 byte)
byte switchRaw = buf.get()
params.switch = switchRaw == 1 ? "on" : "off"
result.params = params
return result
}脚本引擎执行
@Component
public class ScriptEngineExecutor {
private final Map<String, CompiledScript> scriptCache = new ConcurrentHashMap<>();
public JsonObject execute(String productKey, byte[] rawData) {
String scriptContent = getActiveScript(productKey);
// 使用缓存编译脚本
CompiledScript script = scriptCache.computeIfAbsent(productKey, k -> {
ScriptEngineManager manager = new ScriptEngineManager();
ScriptEngine engine = manager.getEngineByName("groovy");
return ((Compilable) engine).compile(scriptContent);
});
// 沙箱绑定
Bindings bindings = new SimpleBindings();
bindings.put("rawData", rawData);
// 在 Sandbox 中执行
Object result = AccessController.doPrivileged(
(PrivilegedAction<Object>) () -> {
try {
return script.eval(bindings);
} catch (ScriptException e) {
throw new ParseException("script execute failed", e);
}
},
// 沙箱权限: 只允许基本运行时权限
new java.security.PermissionCollection() {
{
add(new RuntimePermission("accessDeclaredMembers"));
}
}
);
return JsonParser.parseString(result.toString()).getAsJsonObject();
}
}上报频率控制
防止设备大量高频上报导致平台过载,需要从设备端和平台端两个维度进行控制。
设备端上报间隔
// 设备端上报频率控制
public class UploadThrottler {
private final long minIntervalMs; // 最小上报间隔, 默认 5000ms
private long lastUploadTime = 0;
public UploadThrottler(long minIntervalMs) {
this.minIntervalMs = minIntervalMs;
}
public boolean tryUpload() {
long now = System.currentTimeMillis();
if (now - lastUploadTime >= minIntervalMs) {
lastUploadTime = now;
return true;
}
// 丢弃本次上报 (或做采样降级)
return false;
}
// 弹性上报: 根据数据重要性支持优先级上报
public boolean tryUploadWithPriority(int priority) {
long now = System.currentTimeMillis();
long threshold = minIntervalMs / (priority + 1); // priority 0-4
if (now - lastUploadTime >= threshold) {
lastUploadTime = now;
return true;
}
return false;
}
}平台端限流
@Component
public class UploadRateLimiter {
// 每产品每秒最大消息数
private final Map<String, RateLimiter> productLimiters = new ConcurrentHashMap<>();
public boolean allowUpload(String productKey) {
RateLimiter limiter = productLimiters.computeIfAbsent(productKey,
k -> RateLimiter.create(1000)); // 默认 1000 TPS
return limiter.tryAcquire();
}
// 动态调整限流阈值 (根据平台整体负载)
@EventListener
public void onLoadChange(LoadAlertEvent event) {
double factor = event.getFactor(); // 0.1 ~ 1.0
productLimiters.forEach((k, limiter) -> {
double newRate = limiter.getRate() * factor;
limiter.setRate(Math.max(newRate, 10)); // 最低 10 TPS
});
}
}降级策略
iot:
rate-limit:
default-per-product: 1000 # 每产品默认 TPS
default-per-device: 10 # 每设备默认 TPS
strategy:
overflow: # 超限处理策略
- type: drop # 丢弃 (静默丢弃)
- type: delay # 延迟 (放入延迟队列)
queue-capacity: 10000
delay-ms: 1000
- type: sample # 采样 (只保留关键数据点)
sample-rate: 0.1 # 采样率 10%
global-breach: # 全局过载
- type: reject # 拒绝 (返回错误码)
- type: degrade # 降级 (关闭非核心功能)时序数据存储
使用 TDengine 作为时序数据库存储设备上报数据,基于超级表(STable)模型管理同类设备数据。
超级表设计
-- 创建设备数据超级表
CREATE STABLE IF NOT EXISTS iot.device_data (
ts TIMESTAMP NOT NULL, -- 采集时间
temperature FLOAT COMMENT '温度',
humidity FLOAT COMMENT '湿度',
pressure FLOAT COMMENT '压力',
switch TINYINT COMMENT '开关 0/1',
raw_msg_id VARCHAR(64) COMMENT '原始消息 ID'
) TAGS (
product_key VARCHAR(32), -- 产品标识
device_name VARCHAR(64), -- 设备标识
location VARCHAR(128), -- 安装位置
group_id INT -- 设备分组
);
-- 为每个设备创建子表 (自动继承 STable 结构)
CREATE TABLE IF NOT EXISTS iot.device_pk1_dev001
USING iot.device_data TAGS ('pk1', 'dev001', 'factory-A', 1);
CREATE TABLE IF NOT EXISTS iot.device_pk1_dev002
USING iot.device_data TAGS ('pk1', 'dev002', 'warehouse-B', 2);数据写入
@Service
public class TimeSeriesWriter {
@Autowired
private JdbcTemplate tdJdbcTemplate;
public void writeDeviceData(String productKey, String deviceName,
JsonObject params, long timestamp) {
String tableName = "device_" + productKey + "_" + deviceName;
String sql = "INSERT INTO " + tableName
+ " (ts, temperature, humidity, pressure, switch, raw_msg_id) "
+ "VALUES (?, ?, ?, ?, ?, ?)";
tdJdbcTemplate.update(sql,
new Timestamp(timestamp),
getFloat(params, "temperature"),
getFloat(params, "humidity"),
getFloat(params, "pressure"),
getInt(params, "switch"),
getString(params, "raw_msg_id")
);
}
// 批量写入提升吞吐
@Transactional
public void batchWrite(List<DeviceDataPoint> points) {
NamedParameterJdbcTemplate namedJdbc = new NamedParameterJdbcTemplate(tdJdbcTemplate);
String sql = "INSERT INTO :tableName USING iot.device_data TAGS "
+ "(:productKey, :deviceName, :location, :groupId) "
+ "(ts, temperature, humidity, pressure, switch) "
+ "VALUES (:ts, :temp, :hum, :press, :sw)";
SqlParameterSource[] batch = points.stream()
.map(p -> new MapSqlParameterSource()
.addValue("tableName", "device_" + p.getProductKey() + "_" + p.getDeviceName())
.addValue("productKey", p.getProductKey())
.addValue("deviceName", p.getDeviceName())
.addValue("location", p.getLocation())
.addValue("groupId", p.getGroupId())
.addValue("ts", p.getTimestamp())
.addValue("temp", p.getTemperature())
.addValue("hum", p.getHumidity())
.addValue("press", p.getPressure())
.addValue("sw", p.getSwitchStatus())
)
.toArray(SqlParameterSource[]::new);
namedJdbc.batchUpdate(sql, batch);
}
}典型查询
-- 查询某个设备的最近 1 小时数据
SELECT ts, temperature, humidity
FROM iot.device_pk1_dev001
WHERE ts >= NOW() - 1h
ORDER BY ts ASC;
-- 按产品聚合查询 (跨所有子表)
SELECT AVG(temperature), MAX(temperature), MIN(temperature)
FROM iot.device_data
WHERE product_key = 'pk1'
AND ts >= NOW() - 1d;
-- 按标签过滤查询
SELECT device_name, AVG(temperature) as avg_temp
FROM iot.device_data
WHERE location = 'factory-A'
AND ts >= NOW() - 7d
GROUP BY device_name;
-- 降采样查询 (按 5 分钟窗口聚合)
SELECT INTERVAL(ts, 5m) as window_start,
AVG(temperature) as avg_temp,
MAX(humidity) as max_hum
FROM iot.device_pk1_dev001
WHERE ts >= NOW() - 24h
AND ts < NOW()
INTERVAL(5m);命令下发
平台向设备下发指令,支持同步 RPC、异步下行和批量下发三种模式。
同步 RPC
设备在线时,平台通过 MQTT 发布消息并等待设备应答,适用于需要即时响应的场景。
RPC 流程
平台 设备
| |
|-- MQTT /rpc/invoke/+dev1 -->|
| Request ID: R1 |
| Method: setSwitch |
| Params: {"target":"off"} |
| |
|<-- MQTT /rpc/response/R1 ---|
| Code: 200 |
| Data: {"result":"ok"} |
| |
| 同步等待,超时时间 5s |
| (若超时未收到应答则返回超时) |同步 RPC 实现
@Service
public class SyncRpcService {
private static final long DEFAULT_TIMEOUT_MS = 5000;
private final Map<String, CompletableFuture<RpcResponse>> pendingRequests
= new ConcurrentHashMap<>();
@Autowired
private MqttGateway mqttGateway;
public RpcResponse invoke(String productKey, String deviceName,
String method, JsonObject params) {
return invoke(productKey, deviceName, method, params, DEFAULT_TIMEOUT_MS);
}
public RpcResponse invoke(String productKey, String deviceName,
String method, JsonObject params, long timeoutMs) {
String requestId = UUID.randomUUID().toString();
String topic = "/rpc/invoke/" + productKey + "/" + deviceName;
// 构建请求消息
JsonObject request = new JsonObject();
request.addProperty("id", requestId);
request.addProperty("method", method);
request.add("params", params);
// 注册异步回调
CompletableFuture<RpcResponse> future = new CompletableFuture<>();
pendingRequests.put(requestId, future);
try {
// 发送 MQTT 消息
mqttGateway.sendToTopic(topic, request.toString(), 1);
// 同步等待应答
RpcResponse response = future.get(timeoutMs, TimeUnit.MILLISECONDS);
return response;
} catch (TimeoutException e) {
pendingRequests.remove(requestId);
throw new RpcTimeoutException("device no response within " + timeoutMs + "ms");
} catch (Exception e) {
pendingRequests.remove(requestId);
throw new RpcException("rpc invoke failed", e);
}
}
// 设备端应答回调 (MQTT listener)
public void onRpcResponse(String requestId, String payload) {
CompletableFuture<RpcResponse> future = pendingRequests.remove(requestId);
if (future != null) {
RpcResponse response = JsonParser.parseString(payload).getAsJsonObject();
future.complete(response);
}
}
}设备端处理
// 嵌入式设备端 (C 语言伪代码)
void on_rpc_message(const char* topic, const char* payload) {
// 解析请求
cJSON* root = cJSON_Parse(payload);
char* request_id = cJSON_GetObjectItem(root, "id")->valuestring;
char* method = cJSON_GetObjectItem(root, "method")->valuestring;
cJSON* result = NULL;
int code = 200;
if (strcmp(method, "setSwitch") == 0) {
cJSON* params = cJSON_GetObjectItem(root, "params");
char* target = cJSON_GetObjectItem(params, "target")->valuestring;
// 执行开关操作
gpio_write(SWITCH_PIN, strcmp(target, "on") == 0 ? HIGH : LOW);
result = cJSON_CreateObject();
cJSON_AddStringToObject(result, "result", "ok");
} else {
code = 400;
result = cJSON_CreateObject();
cJSON_AddStringToObject(result, "error", "unknown method");
}
// 构建应答消息
cJSON* response = cJSON_CreateObject();
cJSON_AddStringToObject(response, "id", request_id);
cJSON_AddNumberToObject(response, "code", code);
cJSON_AddItemToObject(response, "data", result);
char* resp_str = cJSON_Print(response);
// 发布到应答 topic
char resp_topic[128];
snprintf(resp_topic, sizeof(resp_topic), "/rpc/response/%s", request_id);
mqtt_publish(resp_topic, resp_str, 1);
free(resp_str);
cJSON_Delete(root);
}异步下行
设备不在线时,将命令缓存到设备影子,设备上线后从影子获取 desired 状态并执行。
异步下发流程
应用 平台 设备
| | |
|-- HTTP PUT /shadow/desired| |
| payload: {"switch":"off"}| |
| | |
| 1. 更新影子 desired |
| 2. 检查设备在线状态 |
| | |
| 3. 设备在线? |
| YES -> 立即推送 desired delta |
| NO -> 等待设备上线 |
| | |
| |<-- MQTT connect -----------|
| | |
| 4. 设备上线回调 |
| 5. 推送缓存的 desired delta |
| | |
| |-- delta: {"switch":"off"}->|
| | |
| |<-- 执行完成 ----------------|
| | |
| 6. 更新 reported,清除已执行 desired |异步命令管理
@Entity
@Table(name = "async_command")
public class AsyncCommand {
@Id
private String commandId;
private String productKey;
private String deviceName;
private String method;
@Column(columnDefinition = "JSON")
private String params;
private Integer commandStatus; // 0-待下发 1-已下发 2-已执行 3-超时 4-失败
private Long ttlSeconds;
private Date gmtCreate;
private Date gmtModified;
}CREATE TABLE async_command (
command_id VARCHAR(64) PRIMARY KEY,
product_key VARCHAR(32) NOT NULL,
device_name VARCHAR(64) NOT NULL,
method VARCHAR(64) NOT NULL,
params JSON NOT NULL,
command_status TINYINT NOT NULL DEFAULT 0 COMMENT '0-待下发 1-已下发 2-已执行 3-超时',
ttl_seconds INT NOT NULL DEFAULT 86400,
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
gmt_modified DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_device_status (product_key, device_name, command_status)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='异步命令表';异步下发调度
@Component
public class AsyncCommandDispatcher {
@Autowired
private AsyncCommandMapper commandMapper;
@Autowired
private ShadowCacheManager shadowCacheManager;
@Autowired
private MqttGateway mqttGateway;
@Transactional
public String dispatchAsync(String productKey, String deviceName,
String method, JsonObject params, long ttlSeconds) {
String commandId = UUID.randomUUID().toString();
AsyncCommand cmd = new AsyncCommand();
cmd.setCommandId(commandId);
cmd.setProductKey(productKey);
cmd.setDeviceName(deviceName);
cmd.setMethod(method);
cmd.setParams(params.toString());
cmd.setCommandStatus(0);
cmd.setTtlSeconds(ttlSeconds);
commandMapper.insert(cmd);
// 检查设备是否在线
if (deviceOnlineChecker.isOnline(productKey, deviceName)) {
// 在线: 立即下发
boolean delivered = deliverToDevice(productKey, deviceName, cmd);
if (delivered) {
cmd.setCommandStatus(1);
commandMapper.updateStatus(commandId, 1);
}
} else {
// 离线: 缓存到影子 desired
shadowCacheManager.mergeDesired(productKey, deviceName,
buildDesiredFromCommand(method, params));
}
return commandId;
}
// 设备上线回调
public void onDeviceOnline(String productKey, String deviceName) {
// 查询待下发的命令
List<AsyncCommand> pending = commandMapper.selectPending(
productKey, deviceName, System.currentTimeMillis());
for (AsyncCommand cmd : pending) {
boolean delivered = deliverToDevice(productKey, deviceName, cmd);
if (delivered) {
cmd.setCommandStatus(1);
commandMapper.updateStatus(cmd.getCommandId(), 1);
}
}
// 同时推送影子 desired
shadowCacheManager.pushDesiredOnOnline(productKey, deviceName);
}
}批量下发
批量任务定义
CREATE TABLE batch_command (
batch_id VARCHAR(64) PRIMARY KEY,
batch_name VARCHAR(128) NOT NULL COMMENT '批量任务名称',
batch_type VARCHAR(32) NOT NULL COMMENT 'upgrade/reboot/config',
method VARCHAR(64) NOT NULL COMMENT '设备端执行方法',
params JSON NOT NULL COMMENT '执行参数',
target_type TINYINT NOT NULL COMMENT '0-按设备列表 1-按产品 2-按分组 3-按区域',
target_value TEXT NOT NULL COMMENT '目标设备列表/产品key/分组ID',
strategy JSON COMMENT '灰度策略',
status TINYINT NOT NULL DEFAULT 0 COMMENT '0-待执行 1-执行中 2-已完成 3-已取消',
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
gmt_modified DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='批量命令任务';批量下发执行
@Service
public class BatchCommandService {
@Autowired
private AsyncCommandDispatcher dispatcher;
@Autowired
private BatchCommandMapper batchMapper;
public String createBatchTask(BatchTaskRequest request) {
String batchId = UUID.randomUUID().toString();
BatchCommand task = new BatchCommand();
task.setBatchId(batchId);
task.setBatchName(request.getBatchName());
task.setBatchType(request.getBatchType());
task.setMethod(request.getMethod());
task.setParams(request.getParams().toString());
task.setTargetType(request.getTargetType());
task.setTargetValue(JsonParser.parseString(request.getTargetValue()));
task.setStrategy(request.getStrategy().toString());
task.setStatus(0);
batchMapper.insert(task);
// 异步执行批量下发
CompletableFuture.runAsync(() -> executeBatch(batchId));
return batchId;
}
public void executeBatch(String batchId) {
BatchCommand task = batchMapper.selectById(batchId);
List<String> targetDevices = resolveTargetDevices(task);
int successCount = 0;
int failCount = 0;
for (String deviceId : targetDevices) {
try {
String[] parts = deviceId.split("/");
String productKey = parts[0];
String deviceName = parts[1];
dispatcher.dispatchAsync(productKey, deviceName,
task.getMethod(),
JsonParser.parseString(task.getParams()).getAsJsonObject(),
86400);
successCount++;
} catch (Exception e) {
log.error("batch dispatch failed for {}: {}", deviceId, e.getMessage());
failCount++;
}
}
batchMapper.updateResult(batchId, successCount, failCount);
}
private List<String> resolveTargetDevices(BatchCommand task) {
// 根据 targetType 解析具体设备列表
JsonObject target = task.getTargetValue();
switch (task.getTargetType()) {
case 0: // 设备列表
return JsonArrayToList(target.getAsJsonArray("devices"));
case 1: // 按产品
return deviceService.listDevicesByProduct(
target.get("productKey").getAsString());
case 2: // 按分组
return deviceService.listDevicesByGroup(
target.get("groupId").getAsLong());
case 3: // 按区域灰度
return deviceService.listDevicesByRegion(
target.get("region").getAsString(),
target.get("percentage").getAsInt());
default:
throw new IllegalArgumentException("unsupported target type: "
+ task.getTargetType());
}
}
}灰度下发策略
{
"strategy": {
"type": "grayscale",
"stages": [
{
"name": "内测",
"percentage": 5,
"conditions": {
"firmware_ver": ">= 2.0.0",
"region": "华东",
"online_time": ">= 7d"
},
"pause_after_minutes": 60,
"observation_window": "30m"
},
{
"name": "小批量",
"percentage": 20,
"pause_after_minutes": 120
},
{
"name": "全量",
"percentage": 100
}
],
"rollback": {
"auto_rollback": true,
"conditions": {
"error_rate": "> 5%",
"offline_rate": "> 10%"
}
}
}
}命令超时与重试机制
超时管理
@Component
public class CommandTimeoutManager {
@Autowired
private AsyncCommandMapper commandMapper;
@Scheduled(fixedRate = 30000) // 每 30 秒扫描一次
public void checkTimeout() {
// 查询所有超时未完成的命令
List<AsyncCommand> timeoutCommands = commandMapper.selectTimeoutCommands(
System.currentTimeMillis());
for (AsyncCommand cmd : timeoutCommands) {
// 超时处理
if (cmd.getRetryCount() < maxRetries) {
// 重试
retryCommand(cmd);
} else {
// 标记为超时终态
commandMapper.updateStatus(cmd.getCommandId(), 3); // 3-超时
// 触发超时回调
commandCallbackService.onCommandTimeout(cmd);
}
}
}
private void retryCommand(AsyncCommand cmd) {
// 指数退避重试
long nextRetryTime = cmd.getNextRetryTime();
if (System.currentTimeMillis() >= nextRetryTime) {
commandMapper.incrementRetry(cmd.getCommandId(),
System.currentTimeMillis() + calculateBackoff(cmd.getRetryCount()));
// 重新入队下发
commandDispatcher.redeliver(cmd);
}
}
private long calculateBackoff(int retryCount) {
// 指数退避: 30s, 1min, 2min, 4min, 8min
return (long) (30_000 * Math.pow(2, retryCount));
}
}重试策略配置
iot:
command:
retry:
max-retries: 5 # 最大重试次数
backoff: exponential # 指数退避
initial-interval-ms: 30000 # 初始间隔
max-interval-ms: 300000 # 最大间隔 (5分钟)
jitter: true # 抖动 (防止重试风暴)
timeout:
sync-rpc-ms: 5000 # 同步 RPC 超时
async-ttl-seconds: 86400 # 异步命令 TTLOTA 升级
OTA(Over-the-Air)升级支持远程固件更新,涵盖固件管理、升级策略、升级流程和断点续传。
固件管理
固件版本管理
CREATE TABLE firmware (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
product_key VARCHAR(32) NOT NULL,
firmware_name VARCHAR(128) NOT NULL,
version VARCHAR(32) NOT NULL COMMENT '语义化版本号',
sign_algorithm VARCHAR(16) NOT NULL DEFAULT 'SHA256' COMMENT '签名算法',
signature VARCHAR(128) NOT NULL COMMENT '固件签名',
file_size BIGINT NOT NULL COMMENT '文件大小 (bytes)',
file_url VARCHAR(512) NOT NULL COMMENT '下载地址 (OSS/CDN)',
file_md5 VARCHAR(64) NOT NULL COMMENT 'MD5 校验值',
diff_patch_url VARCHAR(512) COMMENT '差分升级包地址',
changelog TEXT COMMENT '更新日志',
status TINYINT NOT NULL DEFAULT 0 COMMENT '0-待发布 1-已发布 2-已下架',
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
gmt_modified DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_product_version (product_key, version)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='固件版本';差分升级
差分升级通过对比新旧固件二进制差异,只下发差异部分,减少设备下载量。
# 生成差分补丁 (使用 bsdiff)
bsdiff old_firmware.bin new_firmware.bin patch.bin
# 应用差分补丁
bspatch old_firmware.bin upgrade_firmware.bin patch.bin// 差分升级包管理
public class DiffPatchService {
public DiffPatchInfo generatePatch(String oldVersion, String newVersion,
byte[] oldFirmware, byte[] newFirmware) {
// 调用 bsdiff 生成差分补丁
byte[] patchData = BsDiff.diff(oldFirmware, newFirmware);
DiffPatchInfo info = new DiffPatchInfo();
info.setOldVersion(oldVersion);
info.setNewVersion(newVersion);
info.setPatchSize(patchData.length);
info.setPatchMd5(DigestUtils.md5Hex(patchData));
info.setCompressionRatio(
(double) patchData.length / newFirmware.length);
return info;
}
public byte[] applyPatch(byte[] oldFirmware, byte[] patchData) {
// 设备端应用差分补丁
return BsPatch.patch(oldFirmware, patchData);
}
}升级策略
按设备 / 按批次 / 按区域灰度升级
@Service
public class OtaStrategyService {
public boolean shouldUpgrade(Device device, OtaTask task) {
// 1. 检查设备是否符合灰度条件
OtaStrategy strategy = task.getStrategy();
// 按区域灰度
if (strategy.getRegionFilter() != null
&& !strategy.getRegionFilter().equals(device.getRegion())) {
return false;
}
// 按百分比灰度
if (strategy.getPercentage() < 100) {
// 基于 deviceName 的哈希分桶,保证同一设备灰度一致性
int bucket = Math.abs(device.getDeviceName().hashCode()) % 100;
if (bucket >= strategy.getPercentage()) {
log.debug("device {} bucket {} >= percentage {}, skip",
device.getDeviceName(), bucket, strategy.getPercentage());
return false;
}
}
// 2. 检查当前固件版本
VersionCompareResult cmp = VersionCompare.compare(
device.getFirmwareVersion(), task.getTargetVersion());
if (cmp != VersionCompareResult.LESS_THAN) {
return false; // 设备版本不低于目标版本
}
// 3. 检查设备在线时长 (过滤刚上线设备)
if (strategy.getMinOnlineDays() > 0) {
long onlineDays = ChronoUnit.DAYS.between(
device.getFirstOnlineTime().toInstant(), Instant.now());
if (onlineDays < strategy.getMinOnlineDays()) {
return false;
}
}
return true;
}
}OTA 升级任务定义
CREATE TABLE ota_task (
task_id VARCHAR(64) PRIMARY KEY,
task_name VARCHAR(128) NOT NULL,
product_key VARCHAR(32) NOT NULL,
firmware_id BIGINT NOT NULL,
target_version VARCHAR(32) NOT NULL,
strategy JSON COMMENT '灰度策略配置',
status TINYINT NOT NULL DEFAULT 0 COMMENT '0-待开始 1-执行中 2-已完成 3-已暂停 4-已取消',
total_devices INT NOT NULL DEFAULT 0,
success_count INT NOT NULL DEFAULT 0,
fail_count INT NOT NULL DEFAULT 0,
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
gmt_modified DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='OTA 升级任务';
CREATE TABLE ota_task_device (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
task_id VARCHAR(64) NOT NULL,
product_key VARCHAR(32) NOT NULL,
device_name VARCHAR(64) NOT NULL,
old_version VARCHAR(32) COMMENT '升级前版本',
new_version VARCHAR(32) COMMENT '升级后版本',
status TINYINT NOT NULL DEFAULT 0 COMMENT '0-待升级 1-下载中 2-升级中 3-成功 4-失败 5-跳过',
error_code VARCHAR(64) COMMENT '失败错误码',
progress TINYINT COMMENT '下载进度 0-100',
gmt_create DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
gmt_modified DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_task_device (task_id, product_key, device_name),
INDEX idx_status (status)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='OTA 升级设备明细';升级流程
完整的升级流程为:推送 -> 下载 -> 校验 -> 安装 -> 上报结果。
升级时序
平台 设备
| |
|-- 1. MQTT 推送升级通知 -------->|
| topic: /ota/upgrade |
| payload: {firmware_id, |
| version, url, md5, sign} |
| |
|-- 2. 应答确认收到 ------------>| (device_status -> 下载中)
| |
| 3. 设备下载固件 |
| (HTTPS from OSS/CDN) |
| [断点续传, 分片下载] |
| |
| 4. 校验: MD5 + 签名 |
| 5. 安装固件 |
| (双分区交替/解压覆盖) |
| |
|<-- 6. MQTT 上报升级结果 --------|
| topic: /ota/report |
| payload: {status, version, |
| error_code(可选)} |
| |
| 7. 更新设备影子 / DB |平台端升级分发
@Component
public class OtaDispatcher {
@Autowired
private MqttGateway mqttGateway;
@Autowired
private FirmwareMapper firmwareMapper;
@Autowired
private OtaTaskDeviceMapper taskDeviceMapper;
public void dispatchUpgrade(OtaTask task, List<DeviceUpgrade> devices) {
Firmware firmware = firmwareMapper.selectById(task.getFirmwareId());
for (DeviceUpgrade device : devices) {
String topic = "/ota/upgrade/" + device.getProductKey() + "/" + device.getDeviceName();
JsonObject msg = new JsonObject();
msg.addProperty("firmwareId", firmware.getId());
msg.addProperty("version", firmware.getVersion());
msg.addProperty("url", firmware.getFileUrl());
msg.addProperty("md5", firmware.getFileMd5());
msg.addProperty("signature", firmware.getSignature());
msg.addProperty("signAlgorithm", firmware.getSignAlgorithm());
msg.addProperty("fileSize", firmware.getFileSize());
if (firmware.getDiffPatchUrl() != null) {
msg.addProperty("diffPatchUrl", firmware.getDiffPatchUrl());
}
mqttGateway.sendToTopic(topic, msg.toString(), 2); // QoS 2 确保送达
taskDeviceMapper.updateStatus(device.getTaskId(),
device.getProductKey(), device.getDeviceName(), 1); // 下载中
}
}
}设备端升级处理
// 嵌入式设备 OTA 升级 (C 语言伪代码)
// 升级状态回调
void on_ota_upgrade(const char* topic, const char* payload) {
cJSON* root = cJSON_Parse(payload);
// 解析升级信息
ota_info_t info;
strncpy(info.version, cJSON_GetObjectItem(root, "version")->valuestring, 32);
strncpy(info.url, cJSON_GetObjectItem(root, "url")->valuestring, 512);
strncpy(info.md5, cJSON_GetObjectItem(root, "md5")->valuestring, 64);
info.file_size = cJSON_GetObjectItem(root, "fileSize")->valueint;
// 检查是否有差分升级包
cJSON* diff_url = cJSON_GetObjectItem(root, "diffPatchUrl");
if (diff_url) {
// 差分升级: 下载补丁 + 本地合并
info.is_diff = true;
strncpy(info.diff_url, diff_url->valuestring, 512);
}
cJSON_Delete(root);
// 启动 OTA 任务
ota_start(&info);
}
// OTA 主流程
void ota_start(ota_info_t* info) {
// 1. 准备升级分区
flash_switch_partition(UPGRADE_PARTITION);
// 2. 下载固件 (支持断点续传)
http_download_info_t dl = {
.url = info->url,
.file_size = info->file_size,
.md5 = info->md5,
.dest = UPGRADE_ADDR,
.on_progress = ota_progress_callback,
.resume = true
};
http_result_t result = http_download(&dl);
if (result != HTTP_OK) {
ota_report(OTA_STATUS_FAIL, "download_failed");
return;
}
// 3. 校验 MD5
if (!ota_verify_md5(UPGRADE_ADDR, info->md5)) {
ota_report(OTA_STATUS_FAIL, "md5_mismatch");
return;
}
// 4. 校验签名
if (!ota_verify_signature(UPGRADE_ADDR, info->signature)) {
ota_report(OTA_STATUS_FAIL, "signature_invalid");
return;
}
// 5. 校验固件完整性 (CRC 检查)
if (!ota_verify_integrity(UPGRADE_ADDR)) {
ota_report(OTA_STATUS_FAIL, "integrity_check_failed");
return;
}
// 6. 安装 (标记启动分区)
ota_commit();
ota_report(OTA_STATUS_SUCCESS, info->version);
// 7. 重启设备
system_reboot();
}
// 进度回调
void ota_progress_callback(int percent) {
// 上报下载进度到平台
char topic[64];
snprintf(topic, sizeof(topic), "/ota/progress/%s/%s", PRODUCT_KEY, DEVICE_NAME);
cJSON* msg = cJSON_CreateObject();
cJSON_AddNumberToObject(msg, "progress", percent);
cJSON_AddStringToObject(msg, "version", current_ota_version);
char* str = cJSON_Print(msg);
mqtt_publish(topic, str, 1);
free(str);
cJSON_Delete(msg);
}
// 升级结果上报
void ota_report(int status, const char* detail) {
char topic[64];
snprintf(topic, sizeof(topic), "/ota/report/%s/%s", PRODUCT_KEY, DEVICE_NAME);
cJSON* msg = cJSON_CreateObject();
cJSON_AddNumberToObject(msg, "status", status); // 3-成功 4-失败
cJSON_AddStringToObject(msg, "version", current_ota_version);
if (status == OTA_STATUS_FAIL) {
cJSON_AddStringToObject(msg, "error", detail);
}
char* str = cJSON_Print(msg);
mqtt_publish(topic, str, 2);
free(str);
cJSON_Delete(msg);
}断点续传
固件下载可能因网络波动中断,通过分片下载和已下载字节记录实现断点续传。
设备端断点续传
// 断点续传实现
typedef struct {
uint32_t downloaded; // 已下载字节数
uint32_t file_size; // 总大小
uint32_t chunk_size; // 分片大小 (256KB)
char md5[64]; // 目标 MD5
char url[512]; // 下载 URL
} resume_state_t;
resume_state_t g_resume_state;
http_result_t http_download_resume(http_download_info_t* info) {
// 读取本地已下载进度
uint32_t offset = read_download_checkpoint();
// HTTP Range 请求头
char range_header[64];
snprintf(range_header, sizeof(range_header), "bytes=%u-%u",
offset, offset + info->chunk_size - 1);
http_request_t req = {
.url = info->url,
.headers = {"Range": range_header},
.timeout_ms = 30000
};
http_response_t resp = http_client_get(&req);
if (resp.status_code == 206) { // Partial Content
uint32_t received = offset;
// 追加写入闪存
flash_write(UPGRADE_ADDR + received, resp.body, resp.body_len);
received += resp.body_len;
// 保存断点
save_download_checkpoint(received);
// 上报进度
int percent = (received * 100) / info->file_size;
ota_progress_callback(percent);
if (received >= info->file_size) {
return HTTP_OK; // 下载完成
}
return HTTP_CONTINUE; // 继续下一个分片
}
return HTTP_FAILED;
}平台端下载地址预签名
// 生成带时效的固件下载 URL (OSS 预签名)
public String generateFirmwareDownloadUrl(Long firmwareId, String deviceName) {
Firmware firmware = firmwareMapper.selectById(firmwareId);
// OSS 预签名 URL, 有效期 2 小时
URL url = ossClient.generatePresignedUrl(
firmware.getBucketName(),
firmware.getObjectKey(),
Date.from(Instant.now().plus(2, ChronoUnit.HOURS)),
HttpMethod.GET
);
return url.toString();
}