🔄 SubscribeOperate — 订阅轮询引擎
命名空间: Snet.Core.subscription | 继承: CoreUnify<SubscribeOperate, SubscribeData.Basics> | 接口: ISubscribe, IOn, IOff, IGetStatus
SubscribeOperate 是 Snet 的异步轮询引擎。它周期性地调用用户提供的读取函数,按配置检测数据变化,并通过事件系统分发结果。
配置 (SubscribeData.Basics,继承 SubscribeData.SCData)
| 属性 | 类型 | 默认值 | 说明 |
|---|---|---|---|
SN |
string |
自动 GUID | 唯一标识符 |
Address |
Address? |
null |
订阅的地址集合 |
FunctionAsync |
Func<Address, CancellationToken, Task<OperateResult>>? |
null |
核心 — 用户提供的读取函数 |
HandleInterval |
int |
1000 |
轮询间隔 (ms) |
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
ChangeOut 模式对比
| ChangeOut | AllOut | 行为 |
|---|---|---|
true |
false |
仅输出实际变化的地址数据 |
true |
true |
输出变化项 + 未变化项 |
false |
— | 每次轮询都输出全部数据 |
使用示例
using Snet.Core.subscription;
// 1. 创建订阅配置
var basics = new SubscribeData.Basics
{
Address = new Address(new List<AddressDetails> {
new("温度", "40001", DataType.Float),
new("压力", "40003", DataType.Float)
}),
FunctionAsync = async (address, token) => await modbus.ReadAsync(address, token),
HandleInterval = 500, // 每 500ms 轮询一次
ChangeOut = true, // 仅变化时通知
TaskNumber = 3
};
// 2. 创建订阅引擎
var subscribe = new SubscribeOperate(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("湿度", "40005", DataType.Float)));
// 6. 停止
await subscribe.OffAsync();
配套类型
SubscribeSource<T> — 数据源管理器
继承: CoreUnify<SubscribeSource<T>, string>
public T Source { get; set; }
public DateTime UpdateTime { get; set; }
public async Task<OperateResult> SetAsync(T Data, CancellationToken token)
public async Task<OperateResult> GetAsync(CancellationToken token)
SubscribeService<T> — 数据源注册中心
继承: CoreUnify<SubscribeService<T>, string>
public async Task<OperateResult> SetAsync(string SN, T Data, CancellationToken token)
public async Task<OperateResult> GetAsync(string SN, CancellationToken token)
