KstopaIOT.MQTT 2.0.0-beta

This is a prerelease version of KstopaIOT.MQTT.
dotnet add package KstopaIOT.MQTT --version 2.0.0-beta
                    
NuGet\Install-Package KstopaIOT.MQTT -Version 2.0.0-beta
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="KstopaIOT.MQTT" Version="2.0.0-beta" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="KstopaIOT.MQTT" Version="2.0.0-beta" />
                    
Directory.Packages.props
<PackageReference Include="KstopaIOT.MQTT" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add KstopaIOT.MQTT --version 2.0.0-beta
                    
#r "nuget: KstopaIOT.MQTT, 2.0.0-beta"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package KstopaIOT.MQTT@2.0.0-beta
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=KstopaIOT.MQTT&version=2.0.0-beta&prerelease
                    
Install as a Cake Addin
#tool nuget:?package=KstopaIOT.MQTT&version=2.0.0-beta&prerelease
                    
Install as a Cake Tool

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 四条硬约束(写错都是静默失败)

  1. 执行唯一:ConsumeCrossChainHandoff 默认 false 是刻意的 —— 订阅即"本进程参与执行",多个消费者同时打开会重复执行;
  2. kstopa/mes_nr/send/… 不存在:跨链没有 send,发过去回执 0 条、两侧日志无任何记录;
  3. 两个 Success 别混:result.Success 只表示"回执在超时前到达",result.Response.Success 才是业务结果;
  4. 对方不接管就只有超时: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.ManagedClient 4.3.7.1207 / Newtonsoft.Json 13.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→IsSuccess
  • KstopaIOTRpcResponse 补全缺失字段(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 保证消息入队顺序
  • 修复 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 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. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

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