订阅引擎 - Snet Docs

🔄 SubscribeOperate — 订阅轮询引擎

命名空间: Snet.Core.subscription | 继承: CoreUnify<SubscribeOperate, SubscribeData.Basics> | 接口: ISubscribe, IOn, IOff, IStatus

SubscribeOperate 是 Snet 的异步轮询引擎。它周期性地调用用户提供的读取函数,按配置检测数据变化,并通过事件系统分发结果。

配置 (SubscribeData.Basics,继承 SubscribeData.SCData)

属性 类型 默认值 说明
SN string 自动 GUID 唯一标识符
Address Address? null 订阅的地址集合
FunctionAsync Func<Address, CancellationToken, Task<OperateResult>>? null 核心 — 用户提供的读取函数
HandleInterval int 1000 轮询间隔 (ms)。自 26.222.1 起为 virtual——数采驱动可 override 默认值(如 DBData.Basics 默认每 10000 ms = 10 秒采集一次)
ChangeOut bool true 仅数据变化时触发事件 (false=每次轮询都触发)
AllOut bool false ChangeOut=true 时,是否同时输出变化项和未变化项
TaskNumber int 5 并行轮询任务数

API

// 构造函数
public SubscribeOperate(Basics basics)

// 启动/停止轮询
public async Task<OperateResult> OnAsync(CancellationToken token = default)
public async Task<OperateResult> OffAsync(bool hardClose = false, CancellationToken token = default)

// 动态管理订阅地址
public async Task<OperateResult> SubscribeAsync(Address address, CancellationToken token = default)
public async Task<OperateResult> UnSubscribeAsync(Address address, CancellationToken token = default)

// 状态
public async Task<OperateResult> GetStatusAsync(CancellationToken token = default)

轮询流程

OnAsync() 启动
  └── 创建 TaskNumber 个并行任务
        └── 每个任务循环:
              1. Sleep(HandleInterval) ms
              2. 调用 FunctionAsync(Address, token) → 执行读取
              3. 解析结果为 ConcurrentDictionary<string, AddressValue>
              4. 与上次结果比较 (如果 ChangeOut = true)
              5. 触发 OnDataEvent + OnDataEventAsync

自 26.250.1 起,OnAsync / OffAsync 通过内部生命周期信号量串行化;OffAsync 在返回前会取消并等待轮询任务、组包任务及全部队列消费任务结束。当 FunctionAsyncnull 或抛出异常时,本轮中止并通过信息事件上报(自定义订阅轮询异常:...),轮询进入下一轮继续。

ChangeOut 模式对比

ChangeOut AllOut 行为
true false 仅输出实际变化的地址数据
true true 输出变化项 + 未变化项
false 每次轮询都输出全部数据

使用示例

using System;
using System.Collections.Generic;
using Snet.Core.subscription;
using Snet.Model.data;
using Snet.Model.@enum;

// 1. 创建订阅配置
var basics = new SubscribeData.Basics
{
    Address = new Address(new List<AddressDetails> {
        new("温度", "1", DataType.Float),
        new("压力", "3", DataType.Float)
    }),
    FunctionAsync = async (address, token) => await modbus.ReadAsync(address, token),
    HandleInterval = 500,   // 每 500ms 轮询一次
    ChangeOut = true,        // 仅变化时通知
    TaskNumber = 3
};

// 2. 创建订阅引擎
var subscribe = await SubscribeOperate.InstanceAsync(basics);

// 3. 绑定数据事件
subscribe.OnDataEvent += (s, e) => {
    if (e.Status)
        Console.WriteLine($"数据更新: {e.Message}");
};

// 4. 启动轮询
await subscribe.OnAsync();

// 5. 动态添加地址
await subscribe.SubscribeAsync(new Address(new AddressDetails("湿度", "5", DataType.Float)));

// 6. 停止
await subscribe.OffAsync();

配套类型

SubscribeSource<T> — 数据源管理器

继承: CoreUnify<SubscribeSource<T>, string>

public T Source { get; set; }                       // 读写由内部信号量串行化(与 UpdateTime 原子一致)
public DateTime UpdateTime { get; set; }
public async Task<OperateResult> SetAsync(T Data, CancellationToken token = default)   // Data 为 null 时抛 ArgumentNullException
public async Task<OperateResult> GetAsync(CancellationToken token = default)

SetAsync 在锁内更新 UpdateTime,且仅当值实际变化时(按相等性比较,排除 Time 字段)才触发 OnDataEventAsync("Data Update")。

SubscribeService<T> — 数据源注册中心

继承: CoreUnify<SubscribeService<T>, string>

public async Task<OperateResult> SetAsync(string SN, T Data, CancellationToken token = default)   // SN 非空;Data 非 null;释放后抛 ObjectDisposedException
public async Task<OperateResult> GetAsync(string SN, CancellationToken token = default)           // SN 非空(空时抛 ArgumentException)

Dispose / DisposeAsync 幂等(原子交换守卫);并发对同一 SNSetAsync 调用会被合并,不丢失订阅。

参见