在处理海量数据导入场景时,传统的同步读取→聚合→入库模式常常让数据库成为瓶颈。本文结合项目实战,展示如何用 .NET 8 的 System.Threading.Channels 构建一条高性能的生产者-消费者数据管道:一个异步任务负责从数据库逐条读取、聚合去重,另一个异步任务负责批量 BulkCopy 入库,两者通过有界 Channel 并行协作。

场景来源:SwitchData 项目中的语音专线数据同步(VoiceLineService.ParseAsync),单表百万级数据,多字段聚合去重后写入目标表。

一、为什么需要 Channel?

传统单线程流水线:

Read from DB → Aggregate → Insert to DB → Read next → ...

数据库读写是 I/O 密集型操作,当读取和写入串行执行时,磁盘/网络 I/O 无法被充分利用。生产者-消费者模式可以让读取和写入并行,中间用一个线程安全的缓冲队列解耦。

Channel<T> 是 .NET 提供的高性能异步队列,相比 BlockingCollection<T> 有三个关键优势:

特性 BlockingCollection Channel
等待方式 阻塞线程 await 挂起,零线程占用
有界背压 支持 支持,且背压策略更细
消费方式 foreach await foreach

二、核心架构设计

graph LR
    subgraph Producer["生产者 Task A"]
        A1[ExecuteReaderAsync 逐条读取] --> A2[聚合去重 HashSet + string.Join]
        A2 --> A3[WriteAsync 写入 Channel]
    end
    subgraph Channel["Channel T"]
        B[BoundedChannel 容量 = BatchSize x 2]
    end
    subgraph Consumer["消费者 Task B"]
        C1[ReadAllAsync await foreach 消费] --> C2[填充 DataTable BeginLoadData]
        C2 --> C3[BulkCopyAsync 批量入库]
    end
    A3 --> B --> C1
    style Producer fill:#4CAF50,color:#fff
    style Channel fill:#2196F3,color:#fff
    style Consumer fill:#FF9800,color:#fff

三、有界 Channel 配置

using System.Threading.Channels;

int BatchSize = 20000;

// 关键配置:容量是 BatchSize 的 2 倍,让生产者可以先攒够两批数据
BoundedChannelOptions channelOptions = new(BatchSize * 2)
{
    FullMode = BoundedChannelFullMode.Wait,  // 满时等待,提供天然背压
    SingleReader = true,                     // 单消费者,去掉锁竞争
    SingleWriter = true                      // 单生产者,去掉锁竞争
};

Channel<string[]> channel = Channel.CreateBounded<string[]>(channelOptions);

配置要点:

  • BoundedChannelFullMode.Wait:Channel 满时 WriteAsync 会等待,而不是丢弃数据或阻塞。这是最适合生产者-消费者模式的背压策略 — 消费者慢了生产者自然停下,数据库不会被瞬间打爆。
  • SingleReader = SingleWriter = true:项目约定我们明确只有一个生产者和一个消费者,打开这两个开关让 Channel 内部去掉额外的 lock,吞吐量提升明显(Benchmark 显示约 2-3 倍)。

四、生产者:聚合 + 写入

using var reader = await database.ExecuteReaderAsync(command, cancellationToken);

int fieldCount = reader.FieldCount;

// 每个字段一个 HashSet,用来聚合多行的同列值
Dictionary<int, HashSet<string>> values = new(fieldCount);
for (int i = 0; i < fieldCount; i++)
{
    values[i] = new HashSet<string>(StringComparer.OrdinalIgnoreCase);
}

// 记录上一条聚合键,用于检测聚合边界
string lastCity = null, lastNeName = null, lastNumber = null;

try
{
    while (await reader.ReadAsync(cancellationToken))
    {
        string city = reader[cityOrdinal].ToString();
        string neName = reader[neNameOrdinal].ToString();
        string number = reader[numberOrdinal].ToString();

        // 三个聚合键中有任何一个变了,说明一个聚合单元结束
        if (!string.Equals(city, lastCity, StringComparison.OrdinalIgnoreCase) ||
            !string.Equals(neName, lastNeName, StringComparison.OrdinalIgnoreCase) ||
            !string.Equals(number, lastNumber, StringComparison.OrdinalIgnoreCase))
        {
            if (lastCity != null)
            {
                // 将上一个聚合单元输出为一行
                string[] items = new string[fieldCount];
                for (int i = 0; i < fieldCount; i++)
                {
                    values[i].Remove(string.Empty);
                    items[i] = string.Join('|', values[i]);
                }

                // 检测消费者是否已经出错,及时中断
                if (consumerTask.IsCompleted)
                    await consumerTask;  // 异常会在这里抛出

                // 写入 Channel(背压:如果消费者还没消费完一批,这里会 await)
                await channel.Writer.WriteAsync(items, cancellationToken);
            }

            // 开始新聚合单元:清空所有 HashSet
            for (int i = 0; i < fieldCount; i++)
                values[i].Clear();

            lastCity = city; lastNeName = neName; lastNumber = number;
        }

        // 当前行的值加入各自的 HashSet(自动去重)
        for (int i = 0; i < fieldCount; i++)
        {
            values[i].Add(reader[i].ToString());
        }
    }

    // 别忘了最后一个聚合单元
    if (lastCity != null)
    {
        string[] items = new string[fieldCount];
        for (int i = 0; i < fieldCount; i++)
        {
            values[i].Remove(string.Empty);
            items[i] = string.Join('|', values[i]);
        }
        await channel.Writer.WriteAsync(items, cancellationToken);
    }
}
finally
{
    channel.Writer.Complete();  // 告诉消费者:我写完了
}

设计要点:

  1. consumerTask.IsCompleted 轮询:生产者在每次写入前检查消费者状态。一旦消费者 BulkCopy 抛异常,这里会主动 await 捕获并重新抛出,避免继续往 Channel 里塞无效数据。
  2. channel.Writer.Complete() 放在 finally:不管正常结束还是异常退出都要关闭 Writer,否则消费者的 ReadAllAsync 会永远挂起。
  3. 聚合去重用 HashSet:StringComparer.OrdinalIgnoreCase 做不区分大小写的去重,最后用 string.Join(‘|’, …) 拼接成多值字段。

五、消费者:await foreach + Batch BulkCopy

static async Task PersistAsync(
    ChannelReader<string[]> reader,
    string sourceTableName,
    string targetTableName,
    int batchSize,
    CancellationToken cancellationToken)
{
    // 先从源表拿到 Schema,用来构建空 DataTable
    var dataTable = DbHelper.Default.FillSchema(sourceTableName, TableSchemaType.None);

    await database.ClearTableAsync(targetTableName, cancellationToken);

    // await foreach = 自动循环 + 自动等待 + 自动终止
    await foreach (string[] row in reader.ReadAllAsync(cancellationToken))
    {
        dataTable.Rows.Add(row);

        if (dataTable.Rows.Count >= batchSize)
        {
            await database.BulkCopyAsync(
                dataTable,
                targetTableName,
                batchSize: batchSize,
                addColumnMapping: true,
                cancellationToken: cancellationToken);

            dataTable.Clear();
        }
    }

    // 处理尾部不满一批的数据
    if (dataTable.Rows.Count > 0)
    {
        await database.BulkCopyAsync(
            dataTable, targetTableName,
            batchSize: batchSize, addColumnMapping: true,
            cancellationToken: cancellationToken);
    }
}

消费者启动代码:

Task consumerTask = PersistAsync(
    channel.Reader, sourceTableName, targetTableName,
    BatchSize, cancellationToken);

// 生产者循环...
// await consumerTask;  // 最后等它完成

ReadAllAsync 为什么好用?

ChannelReader.ReadAllAsync() 返回 IAsyncEnumerable,配合 await foreach 有三个自动好处:

  1. 自动循环:不需要手动写 while (await reader.WaitToReadAsync())
  2. 自动等待:Channel 没数据时自动 await 挂起,不占用线程池线程
  3. 自动终止:Writer.Complete() 后循环自然结束,不需要手动检查

对比传统写法:

// 传统写法:啰嗦 + 容易漏 Complete 检查
while (await reader.WaitToReadAsync(cancellationToken))
{
    while (reader.TryRead(out var item))
    {
        Process(item);
    }
}

// await foreach:一行搞定
await foreach (var item in reader.ReadAllAsync(cancellationToken))
{
    Process(item);
}

六、BulkCopy 内部的高性能细节

项目里 ExecuteDataTableAsync 使用了一个高性能技巧:

dataTable.BeginLoadData();
try
{
    object[] values = new object[fieldCount];
    while (await reader.ReadAsync(cancellationToken))
    {
        reader.GetValues(values);  // 一次取出整行所有列,避免多次装箱
        dataTable.Rows.Add(values);
    }
}
finally
{
    dataTable.EndLoadData();
}

BeginLoadData 告诉 DataTable:接下来我要批量塞数据,别每塞一行就去触发 Constraints 检查和事件通知。配合 SqlBulkCopy 入库,单批次 20000 行在本地 SQL Server 上耗时仅 200-400ms。

七、CancellationToken 全程传递

所有异步方法签名都带 CancellationToken,从入口一路传到最底层的 SqlCommand:

public static async Task ParseAsync(CancellationToken cancellationToken = default)
{
    await using var reader = await database.ExecuteReaderAsync(command, cancellationToken);
    await channel.Writer.WriteAsync(items, cancellationToken);
    await database.BulkCopyAsync(dataTable, targetTableName, ..., cancellationToken);
}

如果外部(例如 API 请求断开)触发取消:

  1. reader.ReadAsync 会立即取消,抛 OperationCanceledException
  2. channel.Writer.WriteAsync 不会再等待,直接抛出
  3. consumerTask 会自然结束

整个管道干净地退出,不会有悬挂的数据库连接或未完成的 BulkCopy。

八、项目级 PersistAsync 通用封装

数据库层还提供了一个通用版本的 PersistAsync,直接接收 ChannelReader 并批量入库:

// Database.cs 里的通用方法
public async Task PersistAsync<T>(
    ChannelReader<T> reader,
    int batchSize = 10000,
    CancellationToken cancellationToken = default)
    where T : class, new()
{
    var dataTable = DataTableConverter.CreateTable<T>();

    Exception persistException = null;

    await foreach (T record in reader.ReadAllAsync(cancellationToken))
    {
        if (persistException != null)  // 出错后继续消费直到 Channel 关闭
            continue;

        try
        {
            DataTableConverter.AddRow(dataTable, record);

            if (dataTable.Rows.Count >= batchSize)
                await FlushAsync(dataTable, false, batchSize, cancellationToken);
        }
        catch (Exception ex)
        {
            persistException = ex;
        }
    }

    if (persistException != null) throw persistException;

    await FlushAsync(dataTable, true, batchSize, cancellationToken);
}

这个版本做了优雅降级:如果某一批 BulkCopy 失败,它不立刻抛出,而是继续消费剩余的 Channel 数据(快速排空),等 Writer.Complete() 后再统一抛异常。这样可以避免在异常场景下堵塞上游生产者。

九、性能对比

模式 100万行聚合去重入库耗时 数据库 CPU 内存占用
同步单线程 12分15秒 78% 稳定
Parallel.ForEach 多线程写 5分40秒 95% 锁竞争严重 波动大
Channel 管道 3分20秒 62% 稳定

Channel 管道在保持低数据库压力的同时,吞吐是同步模式的 3.6 倍。核心收益来自读写并行BulkCopy 批量写入的组合。

十、注意事项

  1. Channel 容量不要太小:设为 BatchSize 的 1-3 倍较合理。太小会频繁触发背压,消费者还没攒够一批生产者就被迫停了;太大会浪费内存。
  2. 一定要 Complete() Writer:忘记调用会让 ReadAllAsync 永远挂起。用 try/finally 保证。
  3. 异常处理策略要统一:生产者和消费者都要检查对方是否已完成(consumerTask.IsCompleted),避免单边挂死。
  4. await foreach 自动 Dispose:ReadAllAsync 内部拿到的 CancellationToken 会自动传入,不用额外管理。
  5. SqlBulkCopy 的 EnableStreaming = true:超大表可以减少内存占用(本项目单批次 2w 行,暂不需要)。

总结

用 Channel 构建数据管道的核心思想:把数据流拆成生产者(读取+变换)和消费者(批量写入)两个角色,用有界队列解耦它们的节奏差。这套模式在数据导入、ETL、日志处理、实时监控等场景都可复用。

项目里这套 ParseAsync → Channel → PersistAsync → BulkCopy 四件套,已经成为处理百万级表数据同步的标准范式 — 稳定、快速、背压可控。