在处理海量数据导入场景时,传统的同步读取→聚合→入库模式常常让数据库成为瓶颈。本文结合项目实战,展示如何用
.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(); // 告诉消费者:我写完了
}
设计要点:
- consumerTask.IsCompleted 轮询:生产者在每次写入前检查消费者状态。一旦消费者 BulkCopy 抛异常,这里会主动 await 捕获并重新抛出,避免继续往 Channel 里塞无效数据。
- channel.Writer.Complete() 放在 finally:不管正常结束还是异常退出都要关闭 Writer,否则消费者的 ReadAllAsync 会永远挂起。
- 聚合去重用
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
- 自动循环:不需要手动写 while (await reader.WaitToReadAsync())
- 自动等待:Channel 没数据时自动 await 挂起,不占用线程池线程
- 自动终止: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 请求断开)触发取消:
- reader.ReadAsync 会立即取消,抛 OperationCanceledException
- channel.Writer.WriteAsync 不会再等待,直接抛出
- consumerTask 会自然结束
整个管道干净地退出,不会有悬挂的数据库连接或未完成的 BulkCopy。
八、项目级 PersistAsync 通用封装
数据库层还提供了一个通用版本的 PersistAsync
// 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 批量写入的组合。
十、注意事项
- Channel 容量不要太小:设为 BatchSize 的 1-3 倍较合理。太小会频繁触发背压,消费者还没攒够一批生产者就被迫停了;太大会浪费内存。
- 一定要 Complete() Writer:忘记调用会让 ReadAllAsync 永远挂起。用 try/finally 保证。
- 异常处理策略要统一:生产者和消费者都要检查对方是否已完成(consumerTask.IsCompleted),避免单边挂死。
- await foreach 自动 Dispose:ReadAllAsync 内部拿到的 CancellationToken 会自动传入,不用额外管理。
- SqlBulkCopy 的 EnableStreaming = true:超大表可以减少内存占用(本项目单批次 2w 行,暂不需要)。
总结
用 Channel 构建数据管道的核心思想:把数据流拆成生产者(读取+变换)和消费者(批量写入)两个角色,用有界队列解耦它们的节奏差。这套模式在数据导入、ETL、日志处理、实时监控等场景都可复用。
项目里这套 ParseAsync → Channel → PersistAsync → BulkCopy 四件套,已经成为处理百万级表数据同步的标准范式 — 稳定、快速、背压可控。