KstopaIOT.MQTT
2.0.0-beta
dotnet add package KstopaIOT.MQTT --version 2.0.0-beta
NuGet\Install-Package KstopaIOT.MQTT -Version 2.0.0-beta
<PackageReference Include="KstopaIOT.MQTT" Version="2.0.0-beta" />
<PackageVersion Include="KstopaIOT.MQTT" Version="2.0.0-beta" />
<PackageReference Include="KstopaIOT.MQTT" />
paket add KstopaIOT.MQTT --version 2.0.0-beta
#r "nuget: KstopaIOT.MQTT, 2.0.0-beta"
#:package KstopaIOT.MQTT@2.0.0-beta
#addin nuget:?package=KstopaIOT.MQTT&version=2.0.0-beta&prerelease
#tool nuget:?package=KstopaIOT.MQTT&version=2.0.0-beta&prerelease
KstopaIOT.MQTT
Kstopa 物联平台 MQTT 客户端封装库(前身 IKMqttClient v4.7.x,更名后版本线重置为 1.0.0-beta;当前 2.0.0-beta,主题契约与跨链能力见 新增功能)
Kstopa MQTT 客户端封装库 — 连接佰奥物联 Kstopa 平台 MQTT Broker,提供遥测接入、RPC 写操作、强制刷新上传、跨链交办(nr ↔ mes)、异步回调等待,内置每设备独立 Channel 缓冲、统计监控、精细化设备状态感知和完整信息快照。
| NuGet | KstopaIOT.MQTT 2.0.0-beta(测试版) |
| 框架 | .NET 8 |
| 依赖 | MQTTnet.Extensions.ManagedClient 4.3.7.1207 / Newtonsoft.Json 13.0.4 |
| 命名空间 | KstopaIOT.MQTT |
| 测试 | 编译 0 错误 ✅ |
目录
新增功能(v2.0.0-beta)
相对 1.0.0-beta 基线两项:主题契约改族命名空间(破坏性)+ 跨链交办(nr ↔ mes)。 载荷模型一行未改 ——
KstopaIOTRpcRequest/KstopaIOTRpcResponse/KstopaIOTData/KstopaIOTFlushResponse的字段与 JSON 结构照旧,变的只有主题与新增能力。
一、主题契约:族命名空间(破坏性变更)
主题由 kstopa/{类型}/{设备}/{分组}/… 改为 kstopa/{族}/{类型}/{设备}/{分组}/…,族 = mes(平台/网关链)、nr(Node-RED 链)、mes_nr / nr_mes(跨链转交)。本库代表 mes 链,前缀固定 kstopa/mes/。
| 用途 | 本库主题 |
|---|---|
| 写值 | 发 kstopa/mes/send/{设备}/{分组}/{请求Id}(6 段;分组为空时保留空段,如 kstopa/mes/send/PLC1//req-1) |
| 写值回执 | 订 kstopa/mes/response/{设备}/#(response 与 send 平级,不再是 send/response;回执末尾保留尾斜杠) |
| 上行遥测 / 事件 / 触发 | 订 kstopa/mes/{telemetry\|event\|trigger}/{设备}/# |
| 连接态心跳 | 订 kstopa/mes/connectstatus/{设备}/#(没有分组段) |
| 强制刷新 | 发 kstopa/mes/cmd/flush;回执订 kstopa/mes/cmd/flush/response/{设备}/# |
- 入站解析的设备段从第 4 段起:设备
ts[3]、分组ts[4](旧契约是ts[2]/ts[3],段位整体 +1)—— 这两处改错都是静默错位,不报错; - 主题第 3 段是变量分组名(不是子设备名);分组在网关侧的变量分组表里定义,分组的"使用者"决定这条分组的数据与命令归
mes还是nr; - 旧主题不再支持:
kstopa/send/*、kstopa/telemetry/*、kstopa/cmd/flush… 全部不再处理,没有兼容开关,必须与网关同版升级。
二、跨链交办(nr ↔ mes)
本库同时是交办发起方(平台 → NR)与交办执行方(NR → 平台)。两个方向都走 exe 族:平台 → NR 用 mes_nr、NR → 平台用 nr_mes。
send只用于"让网关执行写值",不参与跨链 ——kstopa/mes_nr/send/…不存在;kstopa/nr/send/… + $handoff是已废弃的旧形状。exe族本身就表示"交给对方处理",载荷不带$handoff标记。
2.1 发起方:平台 → Node-RED
| API | 发布主题 | 回执 |
|---|---|---|
PublishHandoffWriteAsync(req) |
kstopa/mes_nr/exe/write/{设备}/{分组}/{Id} |
不等、无 |
SendHandoffWriteAsync(req, timeout) |
同上 | 等,两种回执见下 |
PublishHandoffDataAsync(device, group, subType, data) |
kstopa/mes_nr/exe/{telemetry\|event\|trigger}/{设备}/{分组}/ |
无(数据类按契约不回执);subType 用共享枚举 KstopaSubType(只接受 Telemetry/Event/Trigger) |
var client = new KstopaIOTMqttClient();
await client.StartClientAsync(new KstopaIOTConnectionConfig
{
MqttIp = "127.0.0.1", MqttPort = 1883,
MqttUName = "kstopa", MqttUPwd = "kstopa123",
Devices = new[] { "PLC1" },
SubscribeNodeRedResponse = true, // 默认即 true:不订就等不到交办结果
});
// 交办写值 + 等结果(按同一个请求 Id 关联回执)
var r = await client.SendHandoffWriteAsync(new KstopaIOTRpcRequest
{
Device = "PLC1", Group = "TechParam", Id = "my-1",
Data = new Dictionary<string, object> { ["温度上限"] = 85 }
}, TimeSpan.FromSeconds(5));
Console.WriteLine($"{(r.Success ? "回执已到" : "超时")} 耗时 {r.ElapsedTime.TotalMilliseconds:F0}ms");
Console.WriteLine($"业务结果:{r.Response?.Success} {r.Response?.ResultMessage}");
// 交办数据(数据类:无请求 Id、无回执)→ kstopa/mes_nr/exe/telemetry/PLC1/TechParam/
await client.PublishHandoffDataAsync("PLC1", "TechParam", KstopaSubType.Telemetry,
new Dictionary<string, object> { ["温度"] = 80.5 });
等结果时回执有两种,都会命中(都按同一个请求 Id 关联,都进 OnNodeRedResponse 事件):
| 回执主题 | 含义 |
|---|---|
kstopa/nr/response/{设备}/{分组}/ |
NR 把交办落地重发后,网关执行写值的回执(沿用原 Id,最常见) |
kstopa/mes_nr/response/{设备}/{分组}/ |
NR 的显式回执(Success=false 即否决) |
子类型用共享枚举 KstopaSubType(KstopaIOT.Interface,网关 / 本库 / Node-RED 节点三方同一份定义):
Write / Telemetry / Event / Trigger,线上字面量一律通过 KstopaSubTypeWire.ToWire(...) 取 —— 改成员名不影响线上契约,改字面量才是破坏性变更。收到未识别的子类型(如 …/exe/bogus/…)会丢弃并报 OnError,不会猜成某个已知子类型。
2.2 执行方:Node-RED → 平台
订阅开关默认关,要显式打开(原因见下方约束 1):
var config = new KstopaIOTConnectionConfig
{
/* …连接参数… */
ConsumeCrossChainHandoff = true, // 订 kstopa/nr_mes/exe/#,否则收不到交办
};
client.OnCrossChainRequest += async (_, req) =>
{
// req.SubType = write | telemetry | event | trigger
// req.Device / req.Group / req.RequestId / req.Data / req.Topic
var ok = await HandleAsync(req);
if (req.RequestId != null) // 数据类交办没有请求 Id,按契约不回执
await client.PublishCrossChainReplyAsync(req.RequestId, req.Device, req.Group, ok, "已处理");
};
回执发 kstopa/nr_mes/response/{设备}/{分组}/,载荷 {"Id":"…","Success":true|false,"Message":"…"}。
2.3 背压(入站队列 + 出站闸门)
跨链路径默认带背压,结构对齐 KstopaIOTDeviceProxy(Bounded 队列 + SemaphoreSlim 限流 + 单 Reader + 计数):
| 配置 | 默认 | 作用 |
|---|---|---|
CrossChainQueueCapacity |
1024 |
入站交办队列容量(对齐 Proxy.ChannelCapacity) |
CrossChainConcurrency |
20 |
入站处理并发上限(对齐 Proxy.ParallelMaxDegree) |
CrossChainFullMode |
Wait |
队列满时:Wait = 真背压(接收侧 await 入队 → 收包循环被顶住 → 压力回传);DropNewest = 丢弃该条并计入 Dropped(不阻塞收包) |
MaxPendingHandoffMessages |
20000 |
出站闸门:跨链发布前若底层待发数已达此值就抛 InvalidOperationException 拒绝入队(0 = 关闭)。避免堆到 MaxPendingMessages 后被 MQTTnet 静默丢弃(那样等回执只会拿到超时,无法区分) |
MaxPendingMessages |
100000 |
底层 ManagedMqttClient 待发上限(此前写死,现可配) |
var config = new KstopaIOTConnectionConfig
{
/* …连接参数… */
ConsumeCrossChainHandoff = true,
CrossChainQueueCapacity = 256, // 按"能容忍多少积压"配
CrossChainConcurrency = 8, // 越大越能扛突发,越小越早把压力回传
CrossChainFullMode = CrossChainFullMode.Wait, // 默认真背压
};
// ⚠️ 要用背压就必须用异步钩子(库会 await 它,处理期间一直占着并发槽位)
client.CrossChainHandler = async req =>
{
await HandleAsync(req); // 这段慢 → 队列填满 → 收包循环被顶住
if (req.RequestId != null)
await client.PublishCrossChainReplyAsync(req.RequestId, req.Device, req.Group, true, "已处理");
};
- 同步事件
OnCrossChainRequest仍是既有用法(行为不变):但它无法承载背压 —— 在里面 fire-and-forget(Task.Run)时槽位立即释放,限流形同虚设。二者二选一(设了CrossChainHandler就不再触发同步事件)。 - 用门面(
KstopaIOTMqttGateway)时走处理器:交办按request.Device路由到IKstopaIOTMessageHandler.HandleCrossChainRequestAsync(带默认实现,既有实现类不用改就能编译);门面不暴露跨链事件,与 Embedded 门面同一口径(避免双消费)。处理耗时同样占用并发槽位 —— 慢处理照样压得回发布方。回执用gateway.PublishCrossChainReplyAsync(...),统计用gateway.GetCrossChainStats()或gateway.GetInfo().CrossChainStats。 - ⚠️
Wait档的代价:队列满时该连接上所有主题一起等(含心跳、回执)。若处理逻辑本身要在同一条连接上等回执,可能互相等到超时 —— 那种场景选DropNewest,或把容量/并发放大。 - 观测:
GetCrossChainStats()返回Received / Dropped / Processed / PendingCount / InFlightCount / MaxInFlight / BackpressureCount / BackpressureWaitTime / HandoffRefused / …。
实测口径(单进程内 broker,60 条 QoS2 交办、每条处理 30ms):
| 配置 | 发布端 | 到达 | 结果 |
|---|---|---|---|
| 容量 1024 / 并发 20 + 异步钩子 | 34ms | 60/60,32–110ms | 峰值并发 20,无背压(余量充足) |
容量 5 / 并发 2 + 异步钩子(Wait) |
20ms | 60/60,1–924ms | 背压 26 次 / 等待 817ms,峰值并发 2 —— 压力真的顶住了收包循环 |
容量 5 / 并发 2(DropNewest) |
17ms | 8/60 | 丢弃 52 条(计数准确),收包不受阻 |
| 容量 5 / 并发 2 + 同步事件 fire-and-forget | 11ms | 60/60,11ms | 峰值 1、无背压 —— 既有用法不变,但不受限流 |
2.4 四条硬约束(写错都是静默失败)
- 执行唯一:
ConsumeCrossChainHandoff默认false是刻意的 —— 订阅即"本进程参与执行",多个消费者同时打开会重复执行; kstopa/mes_nr/send/…不存在:跨链没有send,发过去回执 0 条、两侧日志无任何记录;- 两个
Success别混:result.Success只表示"回执在超时前到达",result.Response.Success才是业务结果; - 对方不接管就只有超时:NR 完全没处理这条交办时不会有任何回执。
三、新增配置项
| 属性 | 类型 | 默认 | 说明 |
|---|---|---|---|
SubscribeNodeRedResponse |
bool |
true |
订 kstopa/nr/response/# 与 kstopa/mes_nr/response/#(按 Devices 过滤)—— 交办写值要等到结果必须开 |
ConsumeCrossChainHandoff |
bool |
false |
订 kstopa/nr_mes/exe/# 并触发 OnCrossChainRequest(整族订阅,交办的设备可能不在 Devices 里) |
四、新增类型与事件
// NR → mes 交办请求(OnCrossChainRequest 的事件参数)
public class KstopaIOTCrossChainRequest
{
public KstopaSubType SubType { get; set; } // 子类型枚举(线上字面量 write / telemetry / event / trigger)
public string Device { get; set; } // 目标设备(主题第 5 段)
public string Group { get; set; } // 分组(主题第 6 段;空段归一为 null)
public string RequestId { get; set; } // 请求 Id(主题末段;数据类交办为 null)
public Dictionary<string, object> Data { get; set; } // 载荷
public string Topic { get; set; } // 原始主题(排查用)
}
| 事件 | 参数 | 触发时机 |
|---|---|---|
OnNodeRedResponse |
KstopaIOTRpcResponse |
kstopa/nr/response/# 或 kstopa/mes_nr/response/# 回执 —— 与 mes 族的 OnResponse 分开报,既有消费者不会收到陌生事件 |
OnCrossChainRequest |
KstopaIOTCrossChainRequest |
收到 kstopa/nr_mes/exe/#(需 ConsumeCrossChainHandoff = true) |
五、联调工具
demo/KstopaIOT.XChainVerify(源码交付,不在 NuGet 包内,需网关在跑;默认 127.0.0.1:1883 / kstopa / kstopa123):
dotnet run—— 库能力自测 13 项,不写设备(交办被网关跳过、回执报文自己造);dotnet run -- --e2e—— 端到端两条链路,会写一次演示设备:A) mes 交办 → NR 落地 → 网关执行 → 回执按同一 Id 回来;B) NR 收事件 → 转交 mes → 库收到并回执。演示设备没有可写数据变量,A 链的业务回执会是Success=false(机制仍已验证)。
六、升级要点
- 库与网关必须同版(契约断代,无兼容开关),升级后请核对两处:写值回执主题变
kstopa/mes/response/…、设备/分组段位 +1; - 跨链 API 在连接层
KstopaIOTMqttClient上(含OnNodeRedResponse/OnCrossChainRequest事件);门面KstopaIOTMqttGateway目前不转发跨链 API 与事件,只用门面做遥测/写值的场景不受影响。
5 分钟上手
安装
dotnet add package KstopaIOT.MQTT --version 2.0.0-beta
你的第一个 Handler(复制即用)
using KstopaIOT.MQTT;
// 1. 实现处理器 — 只需关心业务逻辑
public class MyHandler : IKstopaIOTMessageHandler
{
private readonly string _device;
public MyHandler(string deviceCode) => _device = deviceCode;
public Task HandleTelemetryAsync(KstopaIOTData data)
{
// data.Device = "OP010"
// data.Group = "TechParam"
// data.Data = { "温度": 80.5, "压力": 1.2 }
Console.WriteLine($"[{_device}] 遥测: {string.Join(", ", data.Data)}");
return Task.CompletedTask;
}
public Task HandleEventAsync(KstopaIOTData data) => Task.CompletedTask;
public Task HandleTriggerAsync(KstopaIOTData data) => Task.CompletedTask;
public Task HandleDeviceConnectAsync(KstopaIOTData data) => Task.CompletedTask;
public Task HandleResponseAsync(KstopaIOTRpcResponse r) => Task.CompletedTask;
}
// 2. 一行启动
var gateway = new KstopaIOTMqttGateway(
handlerFactory: deviceCode => new MyHandler(deviceCode));
await gateway.StartAsync(new KstopaIOTConnectionConfig
{
MqttIp = "192.168.1.100"
});
Console.WriteLine("已连接,等待设备消息...");
// 3. 向设备写入数据(可选)
await gateway.PublishWriteAsync(new KstopaIOTRpcRequest
{
Device = "OP010",
Group = "TechParam",
Data = new Dictionary<string, object> { ["温度上限"] = 80 }
});
// 4. 关闭
await gateway.DisposeAsync();
运行后,任何设备上报的遥测数据都会自动路由到 MyHandler.HandleTelemetryAsync。
核心概念
消息从 Broker 到你的 Handler,经历了什么?
MQTT Broker 推送消息
│
▼
┌─ KstopaIOTMqttClient ────────────────────────┐
│ 连接 Broker,接收原始 MQTT 消息 │
│ ├─ 解析 topic(5 种类型自动识别) │
│ ├─ JSON 反序列化为 KstopaIOTData / KstopaIOTRpcResponse │
│ └─ 同步触发事件 │
└──────────────────┬────────────────────────┘
│
▼
┌─ KstopaIOTMqttGateway(推荐,自动完成以下步骤)──┐
│ 1. 按 deviceCode 自动创建/复用 Proxy │
│ 2. 消息非阻塞写入 Channel(TryWrite) │
│ 3. Event/Trigger: Unbounded Channel 永不丢 │
│ Telemetry/ConnectStatus/Response: Bounded│
└──────────────────┬────────────────────────┘
│
▼
┌─ KstopaIOTDeviceProxy(每设备一个实例)─────────┐
│ 5 个独立 Channel: │
│ ├─ telemetry → 串行消费(Bounded) │
│ ├─ connectstatus → 串行消费(Bounded) │
│ ├─ response → 串行消费(Bounded) │
│ ├─ event → Unbounded + SemaphoreSlim│
│ │ 限流并行(默认 20 并发) │
│ └─ trigger → Unbounded + SemaphoreSlim│
│ 限流并行(默认 20 并发) │
│ │
│ 消费者从 Channel 读取 → 调用你的 Handler │
└──────────────────┬────────────────────────┘
│
▼
await handler.HandleTelemetryAsync(data)
await handler.HandleEventAsync(data)
await handler.HandleTriggerAsync(data)
...
数据结构
每条 MQTT 消息被解析为以下模型,你只需关注 Data 字段:
// 遥测/事件/触发/连接状态 消息
public class KstopaIOTData
{
public string Device { get; set; } // 设备编号,如 "OP010"
public string Group { get; set; } // 分组名,如 "TechParam"
public Dictionary<string, object> Data { get; set; } // 实际数据
}
// RPC 写请求
public class KstopaIOTRpcRequest
{
public string Device { get; set; } // 目标设备
public string Group { get; set; } // 分组
public string Method { get; set; } // 写入方法(write / writecache)
public string Cache { get; set; } // 缓存名称(writecache 时使用)
public string Id { get; set; } // 留空自动生成
public Dictionary<string, object> Data { get; set; } // 要写入的数据
}
// RPC 响应
public class KstopaIOTRpcResponse
{
[JsonProperty("Success")] public bool Success { get; set; } // 是否成功
[JsonProperty("Message")] public string ResultMessage { get; set; } // 结果描述
[JsonProperty("Device")] public string DeviceName { get; set; } // 设备名称
[JsonProperty("Group")] public string Group { get; set; } // 分组
[JsonProperty("Id")] public string RequestId { get; set; } // 请求唯一 ID
[JsonProperty("WriteDetails")] public List<WriteDetail> WriteDetails { get; set; } // 写入详情
}
// 写入详情记录
public class WriteDetail
{
public string Method { get; set; } // 写入方法(Write / WriteCache)
public string Variable { get; set; } // 变量名
public object Value { get; set; } // 写入值
}
// NR → mes 的跨链交办请求(OnCrossChainRequest 事件参数,v2.0.0-beta 新增)
// 主题 kstopa/nr_mes/exe/{子类型}/{设备}/{分组}/[{请求Id}]
public class KstopaIOTCrossChainRequest
{
public KstopaSubType SubType { get; set; } // 子类型枚举(线上 write / telemetry / event / trigger)
public string Device { get; set; } // 目标设备
public string Group { get; set; } // 分组(空段归一为 null)
public string RequestId { get; set; } // 请求 Id(数据类交办为 null)
public Dictionary<string, object> Data { get; set; } // 载荷
public string Topic { get; set; } // 原始主题(排查用)
}
DeviceStats(v4.6.1 起)
通过 GetDeviceStats(code) 或 GetAllDeviceStats() 获取:
var stats = gateway.GetDeviceStats("OP010");
Console.WriteLine($"设备: {stats.DeviceCode} 状态: {stats.ConnectionStatus} ({stats.IsDeviceConnected})");
Console.WriteLine($"遥测: 接收{stats.TelemetryReceived} 丢弃{stats.TelemetryDropped} => 处理{stats.TelemetryProcessed}");
Console.WriteLine($"事件: 接收{stats.EventReceived} 丢弃{stats.EventDropped} => 处理{stats.EventProcessed}");
Console.WriteLine($"触发: 接收{stats.TriggerReceived} 丢弃{stats.TriggerDropped} => 处理{stats.TriggerProcessed}");
Console.WriteLine($"RPC: 接收{stats.ResponseReceived} 丢弃{stats.ResponseDropped} => 处理{stats.ResponseProcessed}");
Console.WriteLine($"连接: 接收{stats.ConnectStatusReceived} 重连{stats.ReconnectCount} 心跳{stats.HeartbeatCount}");
Console.WriteLine($"积压: 事件{stats.EventPendingCount} 触发{stats.TriggerPendingCount}");
Console.WriteLine($"首次连接: {stats.FirstConnectTime:yyyy-MM-dd HH:mm:ss}");
Console.WriteLine($"状态变更: {stats.LastStatusChangeTime:yyyy-MM-dd HH:mm:ss}");
// 遍历所有已订阅设备
foreach (var s in gateway.GetAllDeviceStats())
{
var icon = s.ConnectionStatus switch { "在线" => "🟢", "断连中" => "🔴", _ => "⚪" };
Console.WriteLine($"{icon} {s.DeviceCode}: T={s.TelemetryProcessed} E={s.EventProcessed} H={s.HeartbeatCount}");
}
// DeviceStats 类型定义
public record DeviceStats
{
string DeviceCode; // 设备编号
bool? IsDeviceConnected; // 连接状态: true在线/false离线/null未知
string ConnectionStatus; // 精细化: 从未连接/首次在线/在线/离线/断连中
long ReconnectCount; // 掉线次数(真→假重置,假→假持续 +1)
long HeartbeatCount; // 心跳次数(假→真重置,在线消息 +1)
DateTime? FirstConnectTime; // 首次连接时间(UTC)
DateTime? LastStatusChangeTime; // 最近状态变更时间(UTC)
// 接收(入队前 Interlocked 递增)
long TelemetryReceived; long EventReceived;
long TriggerReceived; long ResponseReceived;
long ConnectStatusReceived; // 连接状态消息
// 丢弃(TryWrite 失败时递增)
long TelemetryDropped; long EventDropped;
long TriggerDropped; long ConnectStatusDropped;
long ResponseDropped;
// 积压(Event/Trigger in-flight 数)
long EventPendingCount; long TriggerPendingCount;
// 衍生:实际处理量 = 接收 - 丢弃
long TelemetryProcessed; long EventProcessed;
long TriggerProcessed; long ResponseProcessed;
}
GatewayInfo(v4.6.1 起)
一键获取网关完整快照:
var info = gateway.GetInfo();
Console.WriteLine($"Broker: {info.MqttIp}:{info.MqttPort}");
Console.WriteLine($"Client: {info.ClientId} 连接:{(info.IsConnected ? "在线" : "离线")}(#{info.ConnectionCount})");
Console.WriteLine($"代理数: {info.ProxyCount}");
foreach (var ds in info.DeviceStatsList)
{
Console.WriteLine($" {ds.DeviceCode}: T={ds.TelemetryProcessed} E={ds.EventProcessed}");
}
架构分层
| 层 | 类 | 你通常需要关心的 |
|---|---|---|
| 门面 | KstopaIOTMqttGateway |
✅ 一行启动,日常使用 |
| 连接 | KstopaIOTMqttClient |
⬜ 需要手动控制时用 |
| 代理 | KstopaIOTDeviceProxy |
⬜ 高级定制时用 |
| 工厂 | KstopaIOTDeviceProxyFactory |
⬜ 高级定制时用 |
使用场景
场景 A:基础接入(Gateway 一行启动)
using KstopaIOT.MQTT;
// 创建并启动 — 所有消息自动路由
var gateway = new KstopaIOTMqttGateway(
handlerFactory: deviceCode => new MyHandler(deviceCode),
logger: msg => _logger.LogInformation(msg));
gateway.OnConnectionChanged += (s, connected) =>
Console.WriteLine(connected ? "已连接" : "已断开");
await gateway.StartAsync(new KstopaIOTConnectionConfig
{
MqttIp = "192.168.1.100",
MqttPort = 1883,
Devices = new[] { "OP010", "OP020" } // 只订阅这两台设备
});
// 运行直到 Ctrl+C
await Task.Delay(-1);
场景 B:带完整错误处理和监控
var gateway = new KstopaIOTMqttGateway(
handlerFactory: deviceCode => new MyHandler(deviceCode, _dbContext),
logger: msg => _logger.LogDebug(msg));
// 连接状态
gateway.OnConnectionChanged += (s, connected) =>
{
if (!connected)
_logger.LogWarning("MQTT 断开,自动重连中...");
};
// 错误
gateway.OnGatewayError += (s, error) =>
_logger.LogError("网关异常: {Error}", error);
// 强制刷新响应(主动刷新设备后,网关的处理结果从这里回)
gateway.OnFlushResponse += (s, flush) =>
_logger.LogInformation("刷新响应: 成功={Success} 设备={Devices} 消息={Message}",
flush.Success, string.Join(",", flush.DeviceNames ?? new string[0]), flush.Message);
// 定时输出监控指标
_ = Task.Run(async () =>
{
while (gateway.IsConnected)
{
await Task.Delay(TimeSpan.FromMinutes(1));
var dropped = gateway.GetDroppedStats("OP010");
var received = gateway.GetReceivedStats("OP010");
_logger.LogInformation(
"连接次数:{ConnCount} 代理数:{ProxyCount} " +
"接收(遥测:{TRcv} 事件:{ERcv} 触发:{TrRcv} RPC:{RRcv}) " +
"丢弃(遥测:{TDrop} 事件:{EDrop} 触发:{TrDrop})",
gateway.ConnectionCount, gateway.ProxyCount,
received.telemetry, received.@event, received.trigger, received.response,
dropped.telemetry, dropped.@event, dropped.trigger);
}
});
await gateway.StartAsync(new KstopaIOTConnectionConfig { ... });
场景 C:RPC 读写操作
// 发后不管 — 适合不需要确认的场景
await gateway.PublishWriteAsync(new KstopaIOTRpcRequest
{
Device = "OP010",
Group = "设备信息",
Data = new Dictionary<string, object> { ["设备状态"] = "running" }
});
// 等待响应 — 需要确认写入成功时使用
var result = await gateway.SendWriteAsync(new KstopaIOTRpcRequest
{
Device = "OP010",
Group = "TechParam",
Data = new Dictionary<string, object>
{
["LimitL01"] = 10.5,
["LimitH01"] = 34.1
}
}, timeout: TimeSpan.FromSeconds(5));
if (result.Success)
_logger.LogInformation("写入成功,耗时 {Elapsed}ms",
result.ElapsedTime.TotalMilliseconds);
else
_logger.LogWarning("写入失败: {Description}", result.ResultMessage);
// 查看逐变量写入详情
if (result.Response?.WriteDetails is { Count: > 0 } details)
{
foreach (var d in details)
{
_logger.LogInformation(" [{Method}] {Variable} = {Value}",
d.Method, d.Variable, d.Value);
}
}
场景 D:设备下线
// 移除代理,停止消费任务
await gateway.RemoveDeviceAsync("OP010");
场景 E:DI 容器集成
// Startup.cs
services.AddSingleton(sp =>
{
var handlerFactory = (string deviceCode) =>
{
// 每个设备一次:创建 Scoped Handler
using var scope = sp.CreateScope();
var db = scope.ServiceProvider.GetRequiredService<AppDbContext>();
var logger = scope.ServiceProvider.GetRequiredService<ILogger<MyHandler>>();
return new MyHandler(deviceCode, db, logger);
};
return new KstopaIOTMqttGateway(
handlerFactory,
msg => sp.GetRequiredService<ILogger<KstopaIOTMqttGateway>>().LogInformation(msg));
});
// 后台启动
public class MqttBackgroundService : BackgroundService
{
private readonly KstopaIOTMqttGateway _gateway;
private readonly IConfiguration _config;
public MqttBackgroundService(KstopaIOTMqttGateway gateway, IConfiguration config)
{
_gateway = gateway;
_config = config;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
await _gateway.StartAsync(new KstopaIOTConnectionConfig
{
MqttIp = _config["Mqtt:Host"],
MqttPort = int.Parse(_config["Mqtt:Port"])
});
try { await Task.Delay(-1, stoppingToken); }
catch (OperationCanceledException) { }
finally { await _gateway.DisposeAsync(); }
}
}
场景 F:手动控制(不用 Gateway)
var client = new KstopaIOTMqttClient();
client.SetLogger(msg => _logger.LogInformation(msg));
// 手动订阅事件 + 创建 Proxy
var factory = new KstopaIOTDeviceProxyFactory();
client.OnTelemetry += (s, data) =>
{
var proxy = factory.GetOrCreate(data.Device, new MyHandler(data.Device));
proxy.EnqueueTelemetry(data);
};
// OnEvent / OnTrigger / OnDeviceConnect / OnResponse 同样处理
await client.StartClientAsync(new KstopaIOTConnectionConfig { ... });
配置
连接参数
| 属性 | 类型 | 默认 | 说明 |
|---|---|---|---|
ClientId |
string |
"lms" |
Broker 侧标识,空则自动生成 |
AppendTimestampToClientId |
bool |
false |
为 true 时在 ClientId 后追加 yyyyMMddHHmmssfff 时间戳后缀 |
MqttIp |
string |
"127.0.0.1" |
Broker 地址 |
MqttPort |
int |
1883 |
端口,默认 1883 |
MqttUName |
string |
null |
用户名(需要认证时) |
MqttUPwd |
string |
null |
密码 |
Devices |
string[] |
null |
只订阅指定设备,null = 订阅全部 |
SubscribeNodeRedResponse |
bool |
true |
订阅 kstopa/nr/response/#、kstopa/mes_nr/response/#(跨链交办的结果回执),见 新增功能 |
ConsumeCrossChainHandoff |
bool |
false |
订阅 kstopa/nr_mes/exe/# 并触发 OnCrossChainRequest;默认关是为满足"执行唯一",多消费者都开会重复执行 |
CrossChainQueueCapacity |
int |
1024 |
跨链交办入站队列容量(对齐 Proxy.ChannelCapacity) |
CrossChainConcurrency |
int |
20 |
跨链交办处理并发上限(对齐 Proxy.ParallelMaxDegree) |
CrossChainFullMode |
CrossChainFullMode |
Wait |
队列满时:Wait = 真背压(顶住收包循环);DropNewest = 丢弃并计数 |
MaxPendingHandoffMessages |
int |
20000 |
出站闸门:待发数达此值即拒绝发布(抛 InvalidOperationException),0 = 关闭 |
MaxPendingMessages |
int |
100000 |
底层 ManagedMqttClient 待发上限(此前写死) |
Proxy 配置(每设备独立)
// 通过 Factory 的 configure 回调设置
var proxy = factory.GetOrCreate("OP010", handler, p =>
{
p.ChannelCapacity = 2048; // 缓冲容量
p.ParallelMaxDegree = 20; // 并行线程数
});
| 参数 | 默认 | 说明 |
|---|---|---|
ChannelCapacity |
1024 | Telemetry/ConnectStatus/Response 的 BoundedChannel 容量,满时 TryWrite 返回 false |
ParallelMaxDegree |
20 | Event/Trigger 的并行处理上限(SemaphoreSlim 并发数) |
按设备规模推荐
| 设备数 | ChannelCapacity |
ParallelMaxDegree |
内存估算 |
|---|---|---|---|
| 1-10 台 | 2048 | 20 | < 50 MB |
| 10-50 台 | 1024(默认) | 10(默认) | < 250 MB |
| 50-200 台 | 512 | 5 | < 500 MB |
| 200+ 台 | 256 | 3 | 按需调整 |
公式:
总内存 ≈ 设备数 × ChannelCapacity × 5 通道 × 单条消息大小
API 速查
KstopaIOTMqttGateway(推荐)
| 成员 | 说明 |
|---|---|
KstopaIOTMqttGateway(Func<string, IKstopaIOTMessageHandler> handlerFactory, Action<string> logger = null) |
构造函数:按设备码创建处理器的工厂(必填)+ 可选日志回调 |
StartAsync(config) |
连接 + 启动消费者 |
StopAsync() |
断开连接 |
DisposeAsync() |
断开 + 清空所有 Proxy + 释放资源 |
PublishWriteAsync(req) |
写请求,不等待响应 |
PublishForceUploadAsync(string[] deviceNames) |
强制刷新上传(不指定=连接配置中的设备) |
SendWriteAsync(req, timeout) |
写请求 + 等待响应(超时默认 3000ms,返回含 WriteDetails) |
RemoveDeviceAsync(code) |
移除设备代理 |
GetDroppedStats(code) |
获取丢弃统计 (telem, event, trig, conn, resp) |
GetReceivedStats(code) |
获取接收统计 (telem, event, trig, resp) |
GetDeviceStats(code) |
获取单设备完整统计对象 Device 统计(接收+丢弃+积压+衍生) |
GetAllDeviceStats() |
获取全部已订阅设备的 List<DeviceStats> |
GetInfo() |
获取网关完整信息快照 GatewayInfo(连接+配置+设备统计+跨链背压统计) |
PublishHandoffWriteAsync(req) |
交办写值给 NR 链(不等待) |
SendHandoffWriteAsync(req, timeout) |
交办写值给 NR 链 + 等结果 |
PublishHandoffDataAsync(device, group, subType, data) |
交办数据给 NR 链(subType = KstopaSubType 的 Telemetry/Event/Trigger) |
PublishCrossChainReplyAsync(id, device, group, success, msg) |
回复 NR 交办的处理结果 |
GetCrossChainStats() |
跨链背压统计(同 GetInfo().CrossChainStats) |
OnConnectionChanged |
连接状态变更事件 |
OnGatewayError |
错误事件 |
OnFlushResponse |
强制刷新响应事件(KstopaIOTFlushResponse) |
IsConnected |
当前是否连接 |
ProxyCount |
当前代理数量 |
ConnectionCount |
累计 MQTT Broker 连接次数 |
KstopaIOTMqttClient(手动控制)
| 成员 | 说明 |
|---|---|
StartClientAsync(config) |
连接 Broker |
StopClientAsync() |
断开 |
PublishRequestWriteAsync(req) |
写(不等待) |
PublishForceUploadAsync(string[] deviceNames) |
强制刷新上传(不指定=连接配置中的设备) |
SendRequestWriteAsync(req, timeout) |
写 + 等待回调(返回含 WriteDetails) |
WaitForCallbackAsync(id, timeout) |
等待指定 ID 的回调 |
PublishHandoffWriteAsync(req) |
交办写值给 NR 链(发 kstopa/mes_nr/exe/write/…,不等待) |
SendHandoffWriteAsync(req, timeout) |
交办写值 + 等结果(按同一请求 Id 关联 nr/response 或 mes_nr/response 回执) |
PublishHandoffDataAsync(device, group, subType, data) |
交办数据给 NR 链(subType = KstopaSubType 的 Telemetry/Event/Trigger,无回执) |
PublishCrossChainReplyAsync(id, device, group, success, msg) |
回复 NR 的跨链交办(发 kstopa/nr_mes/response/…) |
CrossChainHandler(属性) |
跨链交办的异步处理钩子:库 await 它,处理期间占并发槽位 → 背压才生效(与 OnCrossChainRequest 二选一) |
GetCrossChainStats() |
跨链背压统计 CrossChainStats(积压/丢弃/并发峰值/背压次数与等待时长/出站被拒数) |
SetLogger(action) |
注入日志 |
ClearAllPendingTasks() |
取消所有等待中回调 |
事件(Client)
| 事件 | 参数 | 触发时机 |
|---|---|---|
OnTelemetry |
KstopaIOTData |
遥测到达 |
OnEvent |
KstopaIOTData |
事件到达 |
OnTrigger |
KstopaIOTData |
触发器到达 |
OnDeviceConnect |
KstopaIOTData |
设备上线/下线 |
OnResponse |
KstopaIOTRpcResponse |
RPC 响应(含 WriteDetails) |
OnNodeRedResponse |
KstopaIOTRpcResponse |
NR 链回执(kstopa/nr/response/#、kstopa/mes_nr/response/#),与 OnResponse 分开报 |
OnCrossChainRequest |
KstopaIOTCrossChainRequest |
收到 NR → mes 的跨链交办(kstopa/nr_mes/exe/#,需 ConsumeCrossChainHandoff = true) |
OnFlushResponse |
KstopaIOTFlushResponse |
强制刷新响应 |
OnError |
string |
内部错误 |
OnIotGatewayConnect |
bool |
网关连接状态 |
序列化工具
// 对象 ↔ 字典
KstopaIOTMqttSerializer.ConvertObjectToDictionary(obj);
KstopaIOTMqttSerializer.DictionaryToObjectConverter<T>(dict);
// 定制 JSON 设置
KstopaIOTMqttSerializer.Settings.DateFormatString = "yyyy-MM-dd HH:mm:ss";
数据扩展工具
KstopaIOTMqttDataExtension 提供安全类型转换和流式构建器,消除手写 TryGetValue + 类型判断的样板代码。
遥测数据安全取值
public Task HandleTelemetryAsync(KstopaIOTData data)
{
// 判断 key 是否存在
if (data.ContainsKey("temperature"))
{
var temp = data.GetDouble("temperature");
}
// 直接用 KstopaIOTData 取值,省去 data.Data.TryGetValue
var speed = data.GetInt32("speed"); // 不存在返回 0
var status = data.GetString("status"); // 不存在返回 ""
var online = data.GetBool("is_online"); // 不存在返回 false
var ts = data.GetDateTime("timestamp"); // 不存在返回 DateTime.MinValue
// 也支持带默认值
var name = data.GetString("device_name", "unknown");
var count = data.GetInt32("batch_count", -1);
// 或直接从字典取值
var dict = data.Data;
var pressure = dict.SafeGetDouble("pressure", 101.3);
var alarm = dict.SafeGetBool("alarm");
// 安全判断字典中是否存在 key(dict 为 null 时返回 false)
if (dict.SafeContainsKey("nested")) {
var sub = dict.SafeGetDict("nested");
var x = sub.SafeGetInt32("x");
}
}
基础类型安全转换
object value = "123.45";
var d = value.ToDouble(); // 123.45,失败返回 0
var i = value.ToInt32(-1); // -1("123.45" 不是有效 int)
var s = value.ToSafeString(); // "123.45",null 返回 ""
var b = value.ToBoolean(); // false
// 数组转换:单值 / 数组 / 可枚举 统一处理
var arr1 = ((object)42).ToInt32Array(); // [42]
var arr2 = ((object)new[] { 1, 2, 3 }).ToInt32Array(); // [1, 2, 3]
var arr3 = ((object)new List<object> { "a", "b" }).ToStringArray(); // ["a", "b"]
写请求流式构建
var request = new KstopaIOTRpcRequest()
.SetDevice("OP010")
.SetGroup("plc_params")
.AddData("target_speed", 2000)
.AddData("target_temp", 42.5)
.AddData("auto_mode", true)
.AddDataRange(new Dictionary<string, object> { // 批量添加
["limit_low"] = 10,
["limit_high"] = 100,
});
// 等同于传统的逐属性赋值:
// new KstopaIOTRpcRequest { Device = "OP010", Group = "plc_params", ... }
await gateway.SendWriteAsync(request, TimeSpan.FromSeconds(5));
数据扩展 API 速查
| 方法 | 目标 | 说明 |
|---|---|---|
data.GetSByte(key, default) |
KstopaIOTData |
遥测安全取 sbyte |
data.GetByte(key, default) |
KstopaIOTData |
遥测安全取 byte |
data.GetInt16(key, default) |
KstopaIOTData |
遥测安全取 short |
data.GetUInt16(key, default) |
KstopaIOTData |
遥测安全取 ushort |
data.GetInt32(key, default) |
KstopaIOTData |
遥测安全取 int |
data.GetUInt32(key, default) |
KstopaIOTData |
遥测安全取 uint |
data.GetInt64(key, default) |
KstopaIOTData |
遥测安全取 long |
data.GetUInt64(key, default) |
KstopaIOTData |
遥测安全取 ulong |
data.GetFloat(key, default) |
KstopaIOTData |
遥测安全取 float |
data.GetDouble(key, default) |
KstopaIOTData |
遥测安全取 double |
data.GetDecimal(key, default) |
KstopaIOTData |
遥测安全取 decimal |
data.GetString(key, default) |
KstopaIOTData |
遥测安全取 string |
data.GetBool(key, default) |
KstopaIOTData |
遥测安全取 bool |
data.GetDateTime(key, default) |
KstopaIOTData |
遥测安全取 DateTime |
data.GetGuid(key, default) |
KstopaIOTData |
遥测安全取 Guid |
data.GetXxxArray(key) |
KstopaIOTData |
遥测安全取数组(15 种类型,缺失返回空数组) |
dict.SafeGetXxx(key, default) |
Dictionary<string,object> |
安全取字典值(16 种基础类型) |
dict.SafeGetXxxArray(key) |
Dictionary<string,object> |
安全取字典数组(15 种类型,缺失返回空数组) |
dict.SafeGetDict(key) |
Dictionary<string,object> |
安全取子字典 |
dict.SafeContainsKey(key) |
Dictionary<string,object> |
安全判断 key 是否存在(dict 为 null 返回 false) |
data.ContainsKey(key) |
KstopaIOTData |
判断遥测数据是否包含 key |
obj.ToXxx(default) |
object |
安全类型转换(SByte ~ Decimal, bool, DateTime, string, Guid) |
obj.ToXxxArray() |
object |
安全转为对应类型数组 |
request.SetDevice/SetGroup/... |
KstopaIOTRpcRequest |
流式设置属性,返回自身 |
request.AddData(key, value) |
KstopaIOTRpcRequest |
流式添加键值对 |
request.AddDataRange(dict) |
KstopaIOTRpcRequest |
流式批量添加 |
request.SetData(dict) |
KstopaIOTRpcRequest |
流式替换整个数据字典 |
项目结构
KstopaIOT.MQTT/
├── KstopaIOTMqttGateway.cs 门面层 — 一行启动(推荐入口)
├── KstopaIOTMqttClient.cs 连接层 — 纯 MQTT 连接
├── KstopaIOTDeviceProxy.cs 代理层 — Channel 缓冲 + 消费者
├── KstopaIOTDeviceProxyFactory.cs 工厂层 — Proxy 生命周期
├── KstopaIOTConnectionConfig.cs 连接配置
├── KstopaIOTMqttSerializer.cs 序列化工具
├── KstopaIOTMqttDataExtension.cs 数据扩展工具(安全取值+流式构建)
├── KstopaIOTCallbackWaitResult.cs 异步回调结果
├── CrossChainStats.cs 跨链背压/处理统计(v2.0.0-beta)
├── CrossChainFullMode.cs 跨链队列满时行为枚举(v2.0.0-beta)
├── DeviceStats.cs 设备统计对象(v4.6.1)
├── GatewayInfo.cs 网关信息对象(v4.6.1)
├── IKstopaIOTMessageHandler.cs 处理器接口(含跨链交办回调,带默认实现)
├── IKstopaIOTMqttClient.cs 连接层接口
├── KstopaIOT/
│ ├── KstopaIOTData.cs 消息数据模型
│ ├── KstopaIOTRpcRequest.cs 写请求模型
│ ├── KstopaIOTRpcResponse.cs 响应模型
│ ├── WriteDetail.cs 写入详情模型(v4.7.2)
│ ├── KstopaIOTCrossChainRequest.cs 跨链交办请求模型(v2.0.0-beta)
│ └── KstopaIOTFlushResponse.cs 强制刷新响应模型(v4.7.2)
└── README.md
(跨链联调工具另在仓库 `demo/KstopaIOT.XChainVerify`,不随 NuGet 包分发)
更新日志
下表 v1.0.0 为更名后的新版本基线(当前带 -beta 测试版标识);v4.x 为更名前
IKMqttClient的同代码迭代历史,供追溯参考。
v2.0.0-beta (2026-09-24) — 🧭 主题族命名空间 + 🔗 跨链交办(nr ↔ mes)
- 破坏性:主题前缀由
kstopa/改为kstopa/mes/(族命名空间);写值回执由kstopa/send/response/…改为kstopa/mes/response/…(response与send平级);入站段位设备/分组由ts[2]/ts[3]改为ts[3]/ts[4]。旧主题不再处理,无兼容开关,需与网关同版升级 - 新增跨链交办(详见 新增功能):
PublishHandoffWriteAsync/SendHandoffWriteAsync/PublishHandoffDataAsync/PublishCrossChainReplyAsync - 新增事件
OnNodeRedResponse、OnCrossChainRequest;新增类型KstopaIOTCrossChainRequest - 子类型收成共享枚举
KstopaSubType(KstopaIOT.Interface,网关 / 本库 / Node-RED 节点三方共用):CrossChainRequest.SubType与PublishHandoffDataAsync(..., subType, ...)都由 string 改为枚举;线上字面量write/telemetry/event/trigger不变,未识别的子类型丢弃并报OnError - 新增配置项
SubscribeNodeRedResponse(默认true)、ConsumeCrossChainHandoff(默认false) - 跨链背压:入站 Bounded 队列 + 并发限流(
CrossChainQueueCapacity/CrossChainConcurrency/CrossChainFullMode,默认真背压)、出站闸门MaxPendingHandoffMessages、异步钩子CrossChainHandler、统计GetCrossChainStats();MaxPendingMessages由写死改为可配(详见 新增功能 §2.3) - 载荷模型与既有 API 签名/行为均未变(本次是"主题改名 + 纯加法")
v1.0.0 — KstopaIOT 品牌重命名(基线)
- 由
IKMqttClient重命名为KstopaIOT.MQTT:PackageId/AssemblyName/RootNamespace统一为KstopaIOT.MQTT - 版本线重置:
Version= 1.0.0-beta(NuGet 包版本,SemVer 3 段 + 测试版标识;正式发布时去掉-beta);AssemblyVersion/FileVersion= 1.0.0.0(4 段 CLR/文件身份) - 依赖:
MQTTnet.Extensions.ManagedClient4.3.7.1207 /Newtonsoft.Json13.0.4,框架 .NET 8 - 功能与更名前的 v4.7.10 完全一致,无 API 破坏性变更
v1.0.0 (2026-08-17 补充) — 🚀 门面 KstopaIOTMqttGateway 新增 OnFlushResponse 转发
KstopaIOTMqttGateway新增并订阅OnFlushResponse事件(KstopaIOTFlushResponse):底层客户端早已支持强制刷新响应(当时主题kstopa/cmd/flush/response,v2.0.0-beta 起为kstopa/mes/cmd/flush/response),但门面此前未转发,消费端从门面无法拿到刷新结果- 使用:
gateway.OnFlushResponse += (s, flush) => ...,与客户端事件同签名 - 无版本号变更、无破坏性 API 变更(仅门面新增事件转发)
v4.7.10 (2026-07-16) — ♻️ KstopaIOTFlushResponse 命名对齐(IsSuccess → Success)
字段命名对齐(延续 v4.7.6 的 KstopaIOTCallbackWaitResult 统一):
KstopaIOTFlushResponse.IsSuccess→KstopaIOTFlushResponse.Success,与KstopaIOTRpcResponse.Success/KstopaIOTCallbackWaitResult.Success命名体系保持一致- ⚠️ JSON 线格式变更:原属性带
[JsonProperty("IsSuccess")],改名后序列化为"Success",下游解析字段名需同步调整 - 同步更新
KstopaIOTMqttClient中强制刷新响应日志:不再引用已移除的IsSuccess,改为输出原始Topic/Payload,便于排查 - 无新增 API、无破坏性接口签名变更(仅属性重命名 + 序列化字段名调整)
v4.7.9 (2026-07-14) — 🔖 版本号对齐 IoTGateway v4.7.9 发布
- 版本号由 4.7.6 升至 4.7.9,与 IoTGateway 网关本体发布版本保持一致;无 SDK 功能性变更。
v4.7.6 (2026-07-09) — ♻️ KstopaIOTCallbackWaitResult 属性名对齐(IsSuccess → Success)
字段命名对齐:
KstopaIOTCallbackWaitResult.IsSuccess→KstopaIOTCallbackWaitResult.Success,与KstopaIOTRpcResponse.Success命名保持一致(v4.7.0 曾由Success改为IsSuccess,本次回退)- 同步更新
KstopaIOTMqttClient(赋值与对象初始化器)及TestClient示例(DisplayResult/MainViewModel)中的引用 - JSON 线格式不变:该属性无
[JsonProperty]标注,仅 C# 代码层命名调整,不影响序列化
v4.7.2 (2026-06-29) — 🧩 WriteDetails 对象化 + 字段精简 + RPC 增强
WriteDetails 重构:
- 新增
WriteDetail类(Method/Variable/Value强类型属性) KstopaIOTRpcResponse.WriteDetails由string→List<WriteDetail>,输出原生 JSON 数组,不再字符串逃逸- 服务端
RpcResponse.WriteDetails同步重构为List<WriteDetail>
字段精简:
KstopaIOTRpcResponse移除冗余字段:Method、CacheName(不再需要)IsSuccess→Success(C# 属性名与[JsonProperty("Success")]对齐)
新增模型:
KstopaIOTFlushResponse— 强制刷新指令的响应模型(IsSuccess/DeviceNames/Message/Timestamp)WriteDetail— 写入详情记录(Method/Variable/Value)
新增 API:
KstopaIOTMqttClient.SendRequestWriteAsync(request, timeout)— 写请求 + 等待回调KstopaIOTMqttClient.PublishRequestWriteAsync(request)— 写请求(不等待)KstopaIOTMqttClient.OnFlushResponse— 强制刷新响应事件KstopaIOTMqttGateway.SendWriteAsync(request, timeout)— 写请求 + 等待响应
驱动增强:
- 新增
WriteCacheBatchAsync批量缓存写入,每条写入和 flush 均记录 WriteDetails AppendWriteDetail强类型 List 操作,替代 JSON 字符串拼接
Bug 修复:
- 响应 Topic 恢复为设备+分组粒度(去除 ClientId 段),修复 MQTT 客户端模式下 ClientId 路由超时问题
ResponseRpcAsync移除 ClientId 精确路由(e.ClientId在 Client 模式下为网关自身 ID)
v4.7.0 (2026-06-16) — 📝 字段命名优化 + 🔧 配置增强
字段命名优化:
RpcResponse/KstopaIOTCallbackWaitResult/KstopaIOTRpcResponse字段统一清晰命名:Cache→CacheName、Description→ResultMessage、Device→DeviceName、Id→RequestId、Success→IsSuccessKstopaIOTRpcResponse补全缺失字段(Group/Method/CacheName),与RpcResponse完全对齐- 使用
[JsonProperty]保持 JSON 线格式不变,仅优化 C# 代码层命名
配置增强:
KstopaIOTConnectionConfig新增AppendTimestampToClientId字段(默认false),控制是否在 ClientId 后追加时间戳后缀MqttIp默认值设为127.0.0.1
工具扩展:
KstopaIOTMqttDataExtension新增ContainsKey/SafeContainsKey方法
Bug 修复:
DeviceThread.WriteSingleVariableAsync修复IsSuccess逐变量覆盖问题:循环中后续成功的写入不再覆盖之前的失败状态
v4.6.3 (2026-06-15) — 🐛 ClientId 去重修复 + ⏱ 时间模拟测试
Bug 修复:
KstopaIOTMqttClient修复相同配置ClientId导致多客户端 ID 重复问题:配置了 ClientId 时自动追加毫秒时间戳_{yyyyMMddHHmmssfff},未配置时仍用 Guid 生成
测试验证:
- ✅ 时间模拟测试通过(详见下方 性能 & 压力测试)
v4.6.2 (2026-06-12) — 🔍 Event/Trigger 接收诊断增强
诊断增强:
KstopaIOTMqttClient新增 Event/Trigger 接收Debug.WriteLine计数输出,按设备+分组统计,不写日志文件DispatchThrottledAsync异常日志从ex.Message扩展为完整ex.ToString(),便于定位根因
Bug 修复:
- 无业务逻辑变更
v4.6.1 (2026-06-11) — 📡 精细化设备状态 + 计数去歧义
设备连接状态(6 种):
ConnectionStatus(string):从未连接 / 首次在线 / 在线 / 离线 / 断连中 / 未知ReconnectCount:掉线次数,true→false 重置后 +1(0 = 从未断连,1+ = 断连次数)HeartbeatCount:当前在线会话心跳数(false→true 重置后每次在线 +1)FirstConnectTime/LastStatusChangeTime:时间戳追踪
计数语义修复:
ReconnectCount默认 0 唯一表示「从未断连」,每次 true→false 重置为 0 再 +1 计数
行为变更:
DispatchThrottledAsync移除异常重试:handler 抛异常直接记录日志
Bug 修复:
ChannelCapacity { get; init; }在构造函数后赋值无效,Channel 创建移到Start()
测试覆盖: 130 项(127 通过,3 预存竞态/性能测试)
v4.2.1 (2026-06-10) — 🛡️ Handler 异常自愈(v4.6.1 已移除)
Event/Trigger 消费异常自动重试(v4.4.0 中移除:异常直接记录日志,不再重试)
v4.2.0 (2026-06-10) — 🎯 Event/Trigger 零丢失
重大变更(内部实现,API 完全兼容):
- Event/Trigger 消费引擎重写:
BoundedChannel(DropOldest) + RunParallelConsumer→UnboundedChannel + SemaphoreSlim 限流并行- ✅ 运行时零丢失:Unbounded Channel 保证
TryWrite永远成功 - ✅ 精确并行控制:SemaphoreSlim 限制并发 in-flight 数
- ✅ 天然背压:满时 Reader 阻塞,不回压生产者
- ✅ 出队 FIFO:SingleReader 保证消息入队顺序
- ✅ 运行时零丢失:Unbounded Channel 保证
- 修复 BoundedChannel 计数器失效 Bug:
FullMode从DropOldest→Wait,满时TryWrite返回 false,*Dropped计数准确 - 新增诊断属性:
EventPendingCount/TriggerPendingCount实时反映积压数 - 测试覆盖:42 项全通过,含零丢失验证(100 条突发 0 丢失 / 10,000 条 0 丢弃)
v4.1.0 — 原架构版本
- BoundedChannel(DropOldest, 1024) + 多 Task 并行消费
- 基础遥测、事件、触发、RPC 功能
测试报告
功能测试(编译 0 错误 ✅)
| 分类 | 覆盖 |
|---|---|
| Client | 参数校验、超时、取消、ClearAll、SetLogger、双重 Dispose |
| Proxy | 5 种消息分派、单消费者排序、接收计数、连接状态解析(在线/离线/未知)、数据闭环 |
| Event | 零丢失验证、并行度控制、完成即释放、积压追踪 |
| Drop | 容量溢出丢弃计数增量验证 |
| Factory | 缓存、计数、TryGet、RemoveAsync |
| Gateway | null 保护、初始状态、全部统计接口、配置开放、信息快照 |
| Serializer | 对象↔字典互转 |
| DataExtension | 77 项断言全覆盖:安全取值、类型转换、流式构建、综合场景 |
| WriteDetails | 三层模型统一(RpcResponse → KstopaIOTRpcResponse(Mqtt) → JSON Payload),原生 JSON 数组输出,字段命名对齐 |
性能 & 压力测试
| 场景 | 结果 |
|---|---|
| 单线程 Enqueue | 9.8M msg/s |
| 1000 并发等待 | 12ms |
| 10K 并发等待 | 89ms |
| 10 设备 × 10K 并行 | 入队 19ms,5s 内完成,无崩溃 |
| 500K 持续爆发 | 入队 4.2M msg/s,系统稳定 |
常见问题
Q: Handler 方法里能做耗时操作吗?
可以,方法签名为 Task,await 异步调用即可。对于遥测/连接/响应 Channel(串行消费),下一消息会等当前完成。对于 event/trigger Channel(Unbounded + SemaphoreSlim 限流并行),最多 ParallelMaxDegree 个并发。
Q: Channel 满了会怎样?
Event/Trigger(Unbounded): 永不丢,SemaphoreSlim 背压控制积压。
Telemetry/ConnectStatus/Response(Bounded): TryWrite 返回 false,*Dropped 计数器递增,可通过 GetDroppedStats() 监控。
Q: 设备突然断线怎么办?
KstopaIOTMqttClient 内置 MQTTnet 自动重连(5s 间隔 + 100K 消息缓冲)。设备重新上线后消息继续投递。
Q: 能动态增删设备吗?
能。设备发消息时自动创建 Proxy(GetOrCreate),不需要预处理。下线时调用 RemoveDeviceAsync 释放资源。
Q: 内存会不会无限增长?
不会。Proxy 按需创建,RemoveDeviceAsync 释放。Telemetry/ConnectStatus/Response Channel 容量固定(Bounded);Event/Trigger 用 Unbounded Channel 但 SemaphoreSlim 背压控制积压。
Q: 线程安全吗?
Client、Factory、Proxy 全部线程安全。ConcurrentDictionary + Channel<T> + Interlocked 保证。
Q: 如何关闭?
await gateway.DisposeAsync(); // 一条搞定:断开 + 清 Proxy + 释放
| Product | Versions Compatible and additional computed target framework versions. |
|---|---|
| .NET | net8.0 is compatible. net8.0-android was computed. net8.0-browser was computed. net8.0-ios was computed. net8.0-maccatalyst was computed. net8.0-macos was computed. net8.0-tvos was computed. net8.0-windows was computed. net9.0 was computed. net9.0-android was computed. net9.0-browser was computed. net9.0-ios was computed. net9.0-maccatalyst was computed. net9.0-macos was computed. net9.0-tvos was computed. net9.0-windows was computed. net10.0 was computed. net10.0-android was computed. net10.0-browser was computed. net10.0-ios was computed. net10.0-maccatalyst was computed. net10.0-macos was computed. net10.0-tvos was computed. net10.0-windows was computed. |
-
net8.0
- KstopaIOT.Interface (>= 2.0.0-beta)
- MQTTnet.Extensions.ManagedClient (>= 4.3.7.1207)
- Newtonsoft.Json (>= 13.0.4)
NuGet packages
This package is not used by any NuGet packages.
GitHub repositories
This package is not used by any popular GitHub repositories.
| Version | Downloads | Last Updated |
|---|---|---|
| 2.0.0-beta | 48 | 9/24/2026 |
| 1.1.0-beta | 70 | 9/11/2026 |
| 1.0.9-beta | 91 | 9/2/2026 |
| 1.0.8-beta | 73 | 9/1/2026 |
| 1.0.7-beta | 81 | 9/1/2026 |
| 1.0.6-beta | 71 | 9/1/2026 |
| 1.0.5-beta | 84 | 8/20/2026 |
| 1.0.4-beta | 72 | 8/19/2026 |
| 1.0.2-beta | 72 | 8/19/2026 |
| 1.0.1-beta | 73 | 8/19/2026 |
| 1.0.0-beta | 94 | 8/17/2026 |