Channel 模式 - Snet Docs

🔀 Channel 模式

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


概述

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

  • 线程安全、无锁的单生产者单消费者或多生产者多消费者队列
  • 通过有界容量实现背压控制
  • 使用 ValueTask 实现异步兼容的读/写
  • 通过 ChannelData.Basics.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); // 带超时异步读取
    public ValueTask<OperateResult> WriteAsync(T item, CancellationToken token);  // 异步写入
    public bool TryRead(out T value);                    // 非阻塞出队(零分配高性能版本)
    public OperateResult TryRead();                      // 非阻塞出队(OperateResult 包装版本)
    public void ResetChannel();                          // 同步重置通道
    public Task ResetChannelAsync(CancellationToken token); // 异步重置通道
    public int Count { get; }                            // 通道中尚未被读取的数据量
    public ChannelReader<T>? Reader { get; }              // 直接访问读取端
    public ChannelWriter<T>? Writer { get; }              // 直接访问写入端
}

IsSync 模式

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

IsSync 行为
true 写入阻塞直到项目被消费(同步管道)
false 写入将项目入队并立即返回(异步管道)

使用模式

生产者端

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

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

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

消费者端

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

Snet 内部使用场景

Channel 模式在整个框架中广泛使用:

使用场景 Channel 类型 描述
订阅管道 ChannelOperate<Address> 在分发事件之前缓冲订阅数据
日志输出 ChannelOperate<LogEvent> 将日志发送与日志持久化解耦
通信队列 ChannelOperate<byte[]> 将传出帧序列化到单个传输通道

有界 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 -- 对于每个写入都必须成功的控制命令