物联网泛在接入
EnerOS v0.45.0 完善物联网泛在接入能力,内置 MQTT 5.0 broker、CoAP 服务器、LwM2M 设备管理、传感器网络聚合与边缘设备批量 OTA。海量 IoT 设备可直接接入 EnerOS,无需外挂 IoT 平台。接入能力由 eneros-iot crate 提供。
IoT 接入架构
┌──────────────────────────────────────────────────────────────┐
│ 边缘设备层 │
│ PMU / RTU / IED / 智能电表 / 传感器 / 边缘网关 │
└────────┬───────────────┬───────────────┬─────────────────────┘
│ │ │
┌────────▼──────┐ ┌──────▼──────┐ ┌──────▼──────────────────┐
│ MQTT 5.0 │ │ CoAP │ │ LwM2M │
│ Broker │ │ Server │ │ Device Mgmt │
│ (TCP/TLS) │ │ (UDP/DTLS) │ │ (CoAP/DTLS) │
└────────┬──────┘ └──────┬──────┘ └────────┬────────────────┘
│ │ │
┌────────▼───────────────▼─────────────────▼──────────────────┐
│ 协议适配层 │
│ IEC 61850 / IEC 60870 / DNP3 / Modbus / DL/T 698.45 │
│ C37.118 (PMU) / OPC UA │
└────────┬─────────────────────────────────────────────────────┘
│
┌────────▼─────────────────────────────────────────────────────┐
│ 数据处理层 │
│ 边缘聚合 / 滤波 / 时序写入 / OTA 管理 / 设备注册 │
└──────────────────────────────────────────────────────────────┘
| 协议 | 端口 | 适用 | 并发能力 |
|---|---|---|---|
| MQTT 5.0 | 1883 / 8883 | 通用 IoT | 1M 连接 |
| CoAP | 5683 / 5684 | 受限设备 | 100k QPS |
| LwM2M | 5684 | 设备管理 | 100k 设备 |
| IEC 61850 | 102 | IED | - |
| IEC 60870-5-104 | 2404 | RTU | - |
| DNP3 | 20000 | 远动 | - |
| C37.118 | 4712 | PMU | - |
| DL/T 698.45 | 9833 | 智能电表 | - |
MQTT 5.0 Broker
内置高性能 MQTT broker,支持百万级连接、共享订阅、会话持久化:
use eneros_iot::mqtt::{MqttBroker, BrokerConfig, Persistence, QoS};
let broker = MqttBroker::new(BrokerConfig {
bind: "0.0.0.0:1883".parse()?,
tls_bind: Some("0.0.0.0:8883".parse()?),
websocket_bind: Some("0.0.0.0:9001".parse()?),
max_connections: 1_000_000,
max_packet_size: 1 << 20, // 1MB
persistence: Persistence::Sled,
session_expiry: Duration::from_secs(86400),
message_queue_size: 1024,
keep_alive: Duration::from_secs(60),
});
broker.start().await?;
// 订阅主题
broker.subscribe("grid/+/voltage", QoS::AtLeastOnce, |topic, payload| {
let bus_id = topic.split('/').nth(1).unwrap();
let v: f64 = serde_json::from_slice(payload)?;
timeseries.write(bus_id, "voltage", ts, v).await?;
Ok(())
}).await?;
// 发布消息
broker.publish("device/rtu_001/command", QoS::ExactlyOnce, &cmd_payload).await?;
主题约定
grid/{bus_id}/voltage — 电压
grid/{bus_id}/current — 电流
grid/{bus_id}/power — 功率
grid/{bus_id}/frequency — 频率
device/{device_id}/status — 设备状态
device/{device_id}/event — 设备事件
device/{device_id}/command — 设备命令
grid/{feeder_id}/load — 馈线负荷
共享订阅
多个消费者共享一个订阅组,实现负载均衡:
// 订阅组 $share/consumer_group/topic
broker.subscribe("$share/processors/grid/+/voltage", QoS::AtLeastOnce, |topic, payload| {
process_voltage(topic, payload).await
}).await?;
MQTT 配置参数
| 参数 | 默认值 | 说明 |
|---|---|---|
| max_connections | 1,000,000 | 最大连接数 |
| max_packet_size | 1MB | 最大包大小 |
| session_expiry | 86400s | 会话过期时间 |
| message_queue_size | 1024 | 离线消息队列 |
| keep_alive | 60s | 心跳间隔 |
| persistence | Sled | 持久化后端 |
CoAP 服务器
支持受限设备的 CoAP 协议,适用于低功耗、低带宽场景:
use eneros_iot::coap::{CoapServer, Resource, Response, Method};
let server = CoapServer::new("0.0.0.0:5683")?
.dtls("0.0.0.0:5684", dtls_config);
server.resource(Resource::get("voltage", |req| async move {
let bus_id = req.uri_path().last().unwrap_or("default");
let v = read_voltage(bus_id).await?;
Response::content(format!("{:.4}", v))
}));
server.resource(Resource::post("command", |req| async move {
let cmd: Command = serde_json::from_slice(req.payload())?;
handle_command(cmd).await?;
Response::changed()
}));
server.resource(Resource::observe("load", |sub| async move {
// 推送实时负荷
while let Some(load) = load_stream.recv().await {
sub.notify(format!("{:.2}", load)).await?;
}
}));
server.start().await?;
CoAP 方法支持
| 方法 | 说明 | 适用 |
|---|---|---|
| GET | 读取资源 | 数据查询 |
| POST | 创建资源 | 命令下发 |
| PUT | 更新资源 | 配置修改 |
| DELETE | 删除资源 | 设备注销 |
| OBSERVE | 订阅变更 | 实时推送 |
LwM2M 设备管理
支持 OMA LwM2M 设备管理协议,覆盖注册、对象访问、固件更新、固件升级全生命周期:
use eneros_iot::lwm2m::{Lwm2mServer, ObjectId, FirmwareUpdate};
let server = Lwm2mServer::new()
.endpoint("coaps://eneros:5684")
.bootstrap(true)
.registration_lifetime(Duration::from_secs(86400));
// 设备注册回调
server.on_register(|dev| async move {
println!("设备注册: {} ({})", dev.endpoint, dev.objects);
register_in_topology(&dev).await?;
Ok(())
});
server.on_deregister(|dev| async move {
println!("设备注销: {}", dev.endpoint);
deregister_from_topology(&dev).await?;
Ok(())
});
// 读取设备对象
let battery = server.read(device_id, ObjectId::Battery).await?;
println!("电量: {}%", battery.level);
println!("电压: {}mV", battery.voltage);
let location = server.read(device_id, ObjectId::Location).await?;
println!("位置: ({}, {})", location.latitude, location.longitude);
// 触发固件更新
server.write(device_id, ObjectId::FirmwareUpdate,
FirmwareUpdate::download("https://ota.eneros/fw/v2.bin")
.signature(Signature::ecdsa(pubkey))
.apply_mode(ApplyMode::AfterDownload)).await?;
// 监听固件更新进度
server.on_firmware_update_progress(|dev, progress| async move {
println!("设备 {} 固件更新: {:.0}%", dev.endpoint, progress * 100.0);
});
LwM2M 对象
| Object ID | 名称 | 用途 |
|---|---|---|
| 0 | LwM2M Security | 安全配置 |
| 1 | LwM2M Server | 服务器配置 |
| 2 | LwM2M Access Control | 访问控制 |
| 3 | Device | 设备信息 |
| 4 | Connectivity Monitoring | 连接监控 |
| 5 | Firmware Update | 固件更新 |
| 6 | Location | 位置 |
| 7 | Connectivity Statistics | 连接统计 |
| 9 | Software Management | 软件管理 |
| 10 | Cellular Connectivity | 蜂窝连接 |
| 11 | APN Connection | APN 配置 |
传感器网络聚合
支持海量传感器数据的边缘聚合,降低带宽与中心处理压力:
use eneros_iot::sensor::{SensorNetwork, Aggregation, Filter};
use std::time::Duration;
let network = SensorNetwork::new()
.sensor("pmu_1", SensorType::PMU, Duration::from_micros(100))
.sensor("pmu_2", SensorType::PMU, Duration::from_micros(100))
.sensor("feeder_1", SensorType::Feeder, Duration::from_millis(100))
.sensor("transformer_1", SensorType::Transformer, Duration::from_secs(1));
// 边缘聚合
network.aggregate(Aggregation::windowed(
Duration::from_millis(100),
|samples| {
// 中位数滤波(去除异常值)
let mut values: Vec<_> = samples.iter().map(|s| s.value).collect();
values.sort_by(|a, b| a.partial_cmp(b).unwrap());
values.get(values.len() / 2).copied().unwrap_or(0.0)
},
)).await?;
// 多级滤波
network.filter(Filter::chain(vec![
Filter::median(5), // 中位数滤波
Filter::lowpass(0.1), // 低通滤波
Filter::outlier(3.0), // 离群点剔除(3σ)
]));
// 自动写入时序数据库
network.forward_to_timeseries(×eries).await?;
// 异常检测
network.on_anomaly(|sensor, value, baseline| async move {
notify_ops(format!("{} 异常: {} vs baseline {}", sensor, value, baseline)).await
});
聚合策略
| 策略 | 说明 | 适用 |
|---|---|---|
| 均值 | 算术平均 | 平滑数据 |
| 中位数 | 中位数 | 抗异常值 |
| 最大值 | 窗口最大 | 峰值监测 |
| 最小值 | 窗口最小 | 谷值监测 |
| 加权平均 | 按质量加权 | 多源融合 |
| P95 | 95 分位 | 服务水平 |
边缘设备 OTA
批量固件更新,支持差分升级、灰度发布与自动回滚:
use eneros_iot::ota::{OtaManager, OtaCampaign, Strategy, Firmware, Signature, OnFailure};
let ota = OtaManager::new()
.concurrency(50) // 并发升级数
.timeout(Duration::from_secs(300));
let campaign = OtaCampaign::new("firmware-v2.3")
.target("rtu.*") // 匹配所有 RTU
.strategy(Strategy::canary(10) // 10% 灰度
.observe(Duration::from_secs(300))
.promote_on(HealthCheck::all_healthy()))
.firmware(Firmware::from_file("firmware.bin")
.size(2 * 1024 * 1024)) // 2MB
.signature(Signature::ecdsa(private_key))
.delta(Delta::from_previous("v2.2")) // 差分升级
.on_failure(OnFailure::rollback()
.threshold(0.05)) // 失败率 > 5% 全部回滚
.on_progress(|device, progress| async move {
println!("{}: {:.0}%", device, progress * 100.0);
});
let result = ota.run(campaign).await?;
println!("成功: {}/{}", result.success, result.total);
println!("失败: {}", result.failures);
println!("总耗时: {:.1}s", result.elapsed.as_secs_f64());
if result.success < result.total * 90 / 100 {
println!("成功率过低,触发全部回滚");
ota.rollback_all().await?;
}
OTA 策略对比
| 策略 | 速度 | 风险 | 适用 |
|---|---|---|---|
| 全量升级 | 慢 | 中 | 大版本 |
| 差分升级 | 快 | 低 | 补丁版本 |
| 灰度发布 | 中 | 极低 | 生产环境 |
| 滚动升级 | 中 | 低 | 常规更新 |
| 蓝绿部署 | 快 | 中 | A/B 验证 |
设备模型
use eneros_iot::device::{Device, DeviceType, Protocol, GeoPoint};
let device = Device::new("rtu_001")
.ty(DeviceType::RTU)
.protocol(Protocol::IEC61850)
.location(GeoPoint::new(31.23, 121.47))
.tags(vec!["feeder_5", "substation_a"])
.metadata(json!({
"vendor": "西门子",
"model": "SICAM",
"install_date": "2024-03-15",
}));
device.register(&iot).await?;
// 设备状态变化
device.on_status_change(|status| async move {
match status {
Status::Online => println!("设备 {} 上线", device.id),
Status::Offline => notify_ops("设备离线").await,
Status::Faulty => create_work_order(&device).await,
}
});
支持的设备类型
| 类型 | 协议 | 用途 | 采集频率 |
|---|---|---|---|
| PMU | C37.118 | 同步相量 | 100 Hz |
| RTU | IEC 60870 / DNP3 | 远动 | 1 Hz |
| IED | IEC 61850 | 智能电子设备 | 1 Hz |
| 智能电表 | DL/T 698.45 | 计量 | 1/15min |
| 传感器 | MQTT / CoAP | 测量 | 1 Hz |
| 边缘网关 | Modbus / OPC UA | 接入 | - |
| 故障指示器 | LoRaWAN | 馈线故障 | 事件触发 |
| 储能 BMS | CAN | 电池管理 | 10 Hz |
安全
所有 IoT 通信支持加密,按设备类型差异化配置:
use eneros_iot::mqtt::{MqttBroker, TlsConfig, Auth};
let broker = MqttBroker::new(config)
.tls(TlsConfig::mutual()
.client_cert_required(true)
.ca("/etc/eneros/iot-ca.crt")
.cipher_suites(CipherSuite::tls_1_3()))
.auth(Auth::token(jwt_validator)
.claim_check(|claims| async {
// 校验设备证书与 token 绑定
claims.device_id == claims.certificate_cn
}));
// 设备级速率限制
broker.rate_limit(RateLimit::per_device(
100, // 100 msg/s
Duration::from_secs(1),
Action::Throttle(Duration::from_secs(60)),
));
性能指标
| 操作 | 吞吐 / 延迟 | 备注 |
|---|---|---|
| MQTT 并发连接 | 1M | 单节点 |
| MQTT 消息吞吐 | 5M msg/s | 单节点 |
| CoAP 请求 | 100k QPS | < 5ms 延迟 |
| LwM2M 注册 | < 50ms | 单设备 |
| OTA 批量升级(1000 设备) | < 5min | 并发 50 |
| 传感器聚合 | 1M samples/s | 单节点 |
| 设备状态查询 | < 10ms | 单设备 |
| 固件差分包生成 | < 1s | 2MB 固件 |
与其他能力的关系
- 时序原生操作:IoT 数据写入时序库,详见 时序原生操作
- 零信任与安全增强:IoT 通信加密,详见 零信任与安全增强
- 国际化与合规:DL/T 698.45 电能表协议,详见 国际化与合规
- 电网拓扑一等公民:设备注册到拓扑,详见 电网拓扑一等公民
- Agent 智能进阶:异常检测基于 IoT 数据