Channel 模式 - Snet Docs

🔀 Channel 模式

ChannelOperate<T> 封装 System.Threading.Channels,在 Snet 内部提供生产者-消费者管道。


概述

Channel 在 Snet 内部用于将数据生产(设备读取)与数据消费(事件订阅者、日志记录、外部接收器)解耦。Channel 模式提供:

  • 线程安全、无锁的单生产者单消费者或多生产者多消费者队列
  • 通过有界容量实现背压控制
  • 使用 ValueTask 实现异步兼容的读/写
  • 通过 ChannelData.IsSync 配置同步/异步模式

ChannelOperate<T>

public class ChannelOperate<T> : CoreUnify<...>
{
    public bool TryWrite(T item);                        // 非阻塞写入
    public ValueTask<OperateResult> ReadAsync(CancellationToken token);        // 异步读取(返回 OperateResult)
    public ValueTask<OperateResult> ReadWaitAsync(int timeOut, CancellationToken token = default); // 带超时异步读取
    public ValueTask<OperateResult> WriteAsync(T item, CancellationToken token);  // 异步写入
    public bool TryRead(out T value);                    // 非阻塞出队(零分配高性能版本)
    public OperateResult TryRead();                      // 非阻塞出队(OperateResult 包装版本)
    public void ResetChannel();                          // 同步重置通道(已释放则抛 ObjectDisposedException)
    public Task ResetChannelAsync(CancellationToken token = default); // 异步重置通道(同上)
    public int Count { get; }                            // 通道中尚未被读取的数据量
    public ChannelReader<T>? Reader { get; }              // 直接访问读取端
    public ChannelWriter<T>? Writer { get; }              // 直接访问写入端
    public bool IsDisposed { get; }                       // 释放后为 true(不可恢复)
}

IsSync 模式

IsSyncChannelData 配置类上的属性(不在 ChannelOperate<T> 上)。

IsSync 行为
true 手动模式:不启动后台读取任务,由调用方通过 ReadAsync / TryRead 读取数据
false 自动模式:构造时自动启动后台读取任务,通过 OnDataEvent 分发数据;此模式下 ReadAsync / ReadWaitAsync 会抛出异常(CustomException)—— TryRead 仍可用

使用模式

生产者端

var channel = await ChannelOperate<MyData>.InstanceAsync();

// 非阻塞写入(尽力而为)
if (!channel.TryWrite(myData))
{
    // Channel 已满 -- 应用背压或丢弃
}

// 或带背压的异步写入
await channel.WriteAsync(myData, CancellationToken.None);

消费者端

// 连续读取循环
while (!token.IsCancellationRequested)
{
    var item = await channel.ReadAsync(token);
    await ProcessItem(item);
}

Snet 内部使用场景

ChannelOperate<T> 是供应用层使用的便捷封装——框架自身的订阅管道、日志输出与通信队列并不经过它。SubscribeOperate 直接使用 System.Threading.Channels.Channel.CreateBounded(...)


有界 vs 无界

Channel 可以创建时带有容量限制:

// 有界:满时应用背压
var boundedChannel = Channel.CreateBounded<MyData>(100);

// 无界:无限制增长(存在内存耗尽风险)
var unboundedChannel = Channel.CreateUnbounded<MyData>();

对于高吞吐量设备的生产使用,建议使用带有适当容量的有界 Channel。这可以防止在网络中断或消费者停滞期间发生内存耗尽。


最佳实践

  1. 谨慎选择容量 -- 太小会导致生产者停滞;太大浪费内存
  2. 使用 CancellationToken -- 将令牌传递给 ReadAsync/WriteAsync 以进行干净关闭
  3. 监控 Channel 指标 -- 通过 TryWrite 返回值跟踪 Count 和丢弃的项目
  4. 传感器数据优先使用 TryWrite -- 对于传感器数据,丢弃旧样本比阻塞更好
  5. 控制命令优先使用 WriteAsync -- 对于每个写入都必须成功的控制命令