跳到主内容

物联网泛在接入

核心能力

物联网泛在接入

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.01883 / 8883通用 IoT1M 连接
CoAP5683 / 5684受限设备100k QPS
LwM2M5684设备管理100k 设备
IEC 61850102IED-
IEC 60870-5-1042404RTU-
DNP320000远动-
C37.1184712PMU-
DL/T 698.459833智能电表-

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_connections1,000,000最大连接数
max_packet_size1MB最大包大小
session_expiry86400s会话过期时间
message_queue_size1024离线消息队列
keep_alive60s心跳间隔
persistenceSled持久化后端

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名称用途
0LwM2M Security安全配置
1LwM2M Server服务器配置
2LwM2M Access Control访问控制
3Device设备信息
4Connectivity Monitoring连接监控
5Firmware Update固件更新
6Location位置
7Connectivity Statistics连接统计
9Software Management软件管理
10Cellular Connectivity蜂窝连接
11APN ConnectionAPN 配置

传感器网络聚合

支持海量传感器数据的边缘聚合,降低带宽与中心处理压力:

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(&timeseries).await?;

// 异常检测
network.on_anomaly(|sensor, value, baseline| async move {
    notify_ops(format!("{} 异常: {} vs baseline {}", sensor, value, baseline)).await
});

聚合策略

策略说明适用
均值算术平均平滑数据
中位数中位数抗异常值
最大值窗口最大峰值监测
最小值窗口最小谷值监测
加权平均按质量加权多源融合
P9595 分位服务水平

边缘设备 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,
    }
});

支持的设备类型

类型协议用途采集频率
PMUC37.118同步相量100 Hz
RTUIEC 60870 / DNP3远动1 Hz
IEDIEC 61850智能电子设备1 Hz
智能电表DL/T 698.45计量1/15min
传感器MQTT / CoAP测量1 Hz
边缘网关Modbus / OPC UA接入-
故障指示器LoRaWAN馈线故障事件触发
储能 BMSCAN电池管理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单设备
固件差分包生成< 1s2MB 固件

与其他能力的关系

相关文档