接手维护一个老批处理服务之后,我对 IAsyncEnumerable 的态度从“知道有这东西”变成了“离不了它”。那个服务每天凌晨会从数据库拉一张几百万行的订单表做统计,原来的实现是Task<List<Order>>的经典套路:一个ToListAsync()把所有数据先装进内存再说。数据量小时一切正常,到了线上真实数据量,进程 GC 频率高得离谱,CPU 大量花在清理托管堆上。后来把核心读取改成异步流,内存峰值降了一个数量级,那份快乐我现在都记得。这篇文章我不打算把 API 文档念一遍,而是按我自己理解它的顺序讲:先说你到底在解决什么问题,再看接口和编译器的底层设计,然后给几个能直接抄进生产代码的用法,最后说几个特别容易踩的坑。
1. 异步迭代出现的真正理由:别再把整张表“搬”进内存
1.1 传统写法的三个隐患
早些年写数据批处理,我们习惯性先拿一个集合:
List<Order> orders = await _db.Orders .Where(o => o.Status == "Paid") .ToListAsync();这句话背后藏着三个问题。第一,数据库把所有匹配记录一次性全部返回,网络传输、数据读取都是从头等到尾。第二,EF Core 把每一行打成实体对象,再塞进List<Order>,这个 List 在整个方法生命周期甚至更久都躺在内存里。第三,后续任何一次遍历、聚合、筛选,都要先跨过这块已经占好的“大石头”。
我在生产里见过最夸张的例子,是某个统计服务把 800 万条明细ToList()之后再聚合,内存 GC 从几秒一次变成几百毫秒一次,最后整个服务 CPU 卡在垃圾回收上,业务吞吐量直接掉到原来的三分之一。问题不在业务逻辑,而在“数据还没用上,就已经全部常驻内存”。
1.2 Task<IEnumerable > 只是“半截”异步
很多人觉得Task<IEnumerable<T>>已经是异步了,但它的本质只是“等一下再拿集合”。拿到集合之后,你对它做foreach仍然是同步遍历;如果这个集合是惰性的,每一次MoveNext()都可能阻塞调用线程。也就是说,第一段异步只解决了“到达数据前不占线程”的问题,没有解决“拿到数据流之后逐条读取不阻塞”的问题。
IEnumerable<T>本身能流式读取,但它读取的过程是同步的。当数据源真的是网络接口、远程 API、数据库游标这类需要异步 I/O 的地方,你要么让线程傻等,要么手动拆成一页一页去拉。而 IAsyncEnumerable 想解决的问题,就是“异步地从远处拿数据,拿到一条就处理一条,拿不到的时候线程去干别的”。
1.3 三种数据形态放在一起看
| 形态 | 数据产生方式 | 内存表现 | 等待时是否占用线程 |
|---|---|---|---|
IEnumerable<T> | 同步拉取 | 可以惰性,但读取时阻塞 | 占用调用线程 |
Task<IEnumerable<T>> | 异步拿到集合,遍历仍同步 | 通常一次性全量入内存 | 等待期间不占,遍历时占 |
IAsyncEnumerable<T> | 逐个异步拉取 | 边取边甩,天然流式 | 等待期间不占,也不阻塞 |
这个对比基本就是我在选型时的判断依据。数据量只有几十条时,怎么写都无所谓;一旦数据可能是几万、几百上千万条,IAsyncEnumerable 的价值就立刻体现出来。
2. 接口和编译器的三个关键设计:Current、MoveNextAsync 与 ValueTask
2.1 接口本身很短,信息量不小
IAsyncEnumerable 的接口定义非常简洁:
public interface IAsyncEnumerable<out T> { IAsyncEnumerator<T> GetAsyncEnumerator(CancellationToken cancellationToken = default); } public interface IAsyncEnumerator<out T> : IAsyncDisposable { T Current { get; } ValueTask<bool> MoveNextAsync(); }整个抽象就三件事:Current给出当前值,MoveNextAsync()异步判断还有没有下一个,DisposeAsync()负责清理。GetAsyncEnumerator接收一个取消令牌,让外部可以在枚举开始前就挂上取消信号。注意泛型参数带out,意味着IAsyncEnumerable<Dog>可以直接当成IAsyncEnumerable<Animal>用,这一点和IEnumerable<T>的协变语义一致,定义接口时非常有用。
2.2 yield 和 await 混在一起,编译器怎么处理
在 C# 8 里写异步迭代器,最直观的形态长这样:
async IAsyncEnumerable<int> ProduceNumbers(CancellationToken token = default) { for (int i = 0; i < 100; i++) { await Task.Delay(10, token); yield return i; } }第一次看到这种代码的人通常会问:await和yield return怎么能混在一起?关键在于,这不再是普通方法,而是一个异步迭代器。编译器会把它改造成一个状态机:方法被调用时不会立即执行任何业务代码,只是把迭代器“放出去”;每次消费者调用MoveNextAsync(),状态机才跑一段代码,跑到下一个await或者yield return就暂停下来。
可以把它理解成视频播放的“按帧加载”:你要看下一帧,我才去拉下一帧。而ToListAsync()的方式相当于把整季视频一次性下载完再播放。两者最终看到的都是完整数据,但内存曲线和首屏延迟完全不同。
2.3 ValueTask 的选择是性能杀手锏
MoveNextAsync()的返回类型不是Task<bool>,而是ValueTask<bool>,这个细节值得单独说。异步流的特点是:很多情况下下一个元素已经准备好了。比如 Channel 缓冲区里正好有数据,或者包装的同步集合里还有剩余项,这时候MoveNextAsync()完全可以同步返回true,根本不需要一次真正的异步等待。
如果接口设计成Task<bool>,每取一个元素都可能要分配一个 Task 对象。一百万个元素就有一百万次分配,GC 压力相当大。ValueTask的设计思路是“能同步返回就同步返回,需要真正异步时才去等待”,在零分配和高吞吐之间做了很好的平衡。虽然日常使用中你不会感知到它,但这个选择直接关系到百万行数据流的稳定性。
3. await foreach 用法里的几个细节:取消、配置与资源释放
3.1 消费一个异步流的基本逻辑
消费 IAsyncEnumerable 的常规入口是await foreach:
await foreach (var order in GetOrders(ct)) { await ProcessOrder(order, ct); }它背后做的事等价于:
await using (var enumerator = GetOrders(ct).GetAsyncEnumerator(ct)) { while (await enumerator.MoveNextAsync()) { var order = enumerator.Current; await ProcessOrder(order, ct); } }也就是说,循环正常结束、通过break提前退出、或中间抛出异常,DisposeAsync()都会被执行。这个自动清理机制很关键:异步流持有数据库连接、Channel 读取端时,不及时释放会直接造成资源泄漏。
3.2 取消令牌是怎么流进迭代器的
异步迭代器方法里的参数,有一个专门配合取消的属性:
async IAsyncEnumerable<Order> GetOrders( [EnumeratorCancellation] CancellationToken token = default) { while (true) { var page = await _db.Orders.Take(100).ToListAsync(token); foreach (var item in page) { yield return item; } if (page.Count < 100) { break; } } }调用方可以这样挂上令牌:
var cts = new CancellationTokenSource(TimeSpan.FromMinutes(5)); await foreach (var order in GetOrders().WithCancellation(cts.Token)) { await ProcessOrder(order, cts.Token); }WithCancellation这个扩展方法会把令牌传给GetAsyncEnumerator,然后通过[EnumeratorCancellation]特性绑定到迭代器方法内部的token参数上。这样做的价值在于:数据库查询、Task.Delay内部都能感知取消,而不是等下一次MoveNextAsync才抛出异常。
3.3 迭代器里的 try/finally 与资源释放
在异步迭代器内部写try/finally是常见需求,尤其是要保证“即使消费者提前 break,也要把外部资源清理掉”的时候:
async IAsyncEnumerable<Order> GetOrders() { var client = CreateClient(); try { while (await client.NextAsync()) { yield return client.Current; } } finally { await client.CloseAsync(); } }这里finally会在消费者正常遍历完、主动退出、或异常发生时执行。异步迭代器允许在finally里await,这一点和普通方法不同,也是设计者专门为“异步清理”留的口子。
3.4 ConfigureAwait(false) 什么时候加
await foreach也有对应的ConfigureAwait扩展:
await foreach (var item in source.ConfigureAwait(false)) { // ... }我的经验是:在 UI 程序里,默认保留同步上下文可以让异步流里的代码继续回到 UI 线程,更新控件很方便;但在类库、后台服务、中间件里,如果没有回到特定同步上下文的需求,尽量加上ConfigureAwait(false),避免不必要的上下文捕获和切换开销。ASP.NET Core 默认没有同步上下文,问题不大,但养成习惯对公共库更友好。
4. 生产环境里三个高频应用场景:EF Core、Channel 与并行消费
4.1 EF Core 流式读取:AsAsyncEnumerable
EF Core 3.0 以后,IQueryable<T>上可以直接调AsAsyncEnumerable():
await foreach (var order in _db.Orders .Where(o => o.Status == Status.Paid) .AsAsyncEnumerable() .WithCancellation(ct)) { await HandleOrder(order, ct); }这样 EF Core 会走DbDataReader,一次只物化一行实体,整张表不会堆在内存里。代价是:整个await foreach期间,这个 DbContext 一直被占用,不能再同时用同一个上下文发其他查询。如果业务里需要并行做两件事,要么拆两个 DbContext,要么先按小批次缓冲再做别的。
4.2 Channel 拼出生产者-消费者管道
IAsyncEnumerable 和 Channel 是绝配。Channel 本质是线程安全的异步队列,Reader.ReadAllAsync()直接返回IAsyncEnumerable<T>:
var channel = Channel.CreateUnbounded<Order>(new UnboundedChannelOptions { SingleReader = true, SingleWriter = false }); async Task ProduceAsync(CancellationToken ct) { try { await foreach (var order in FetchRemoteOrders(ct)) { await channel.Writer.WriteAsync(order, ct); } } finally { channel.Writer.TryComplete(); } } async Task ConsumeAsync(CancellationToken ct) { await foreach (var order in channel.Reader.ReadAllAsync(ct)) { await ProcessOrder(order, ct); } } var consumeTask = ConsumeAsync(ct); await ProduceAsync(ct); await consumeTask;这个模式的妙处在于天然支持背压。把CreateUnbounded换成CreateBounded并设置容量后,消费者处理不过来时,WriteAsync会自动等待,生产者不会无脑堆积数据。生产者和消费者可以运行在不同的任务上,甚至一个写、多个读,只要把SingleReader设为false即可。
4.3 并行消费:Parallel.ForEachAsync
.NET 6 之后,Parallel.ForEachAsync可以直接吃 IAsyncEnumerable:
await Parallel.ForEachAsync( GetOrders(ct), new ParallelOptions { MaxDegreeOfParallelism = 8, CancellationToken = ct }, async (order, ct) => { await ProcessOrder(order, ct); });这个 API 很适合“每条数据处理都有网络延迟”的场景。比如每条订单都要调用外部价格服务,逐个等和 8 个并发等,总耗时完全不同。要注意:并行消费不保证处理顺序,所以只适用于任务之间没有先后依赖的场景;如果后续逻辑强依赖顺序,老老实实用一个await foreach串行处理,别并行。
5. 我踩过的几个真实坑:重枚举、上下文过期与“假 LINQ”
5.1 一份异步流不能想当然地枚举两次
var stream = GetOrders(ct); await foreach (var order in stream) { /* ... */ } await foreach (var order in stream) { /* ... */ }这个代码最坑的地方是:它不一定报错。对于迭代器方法编写的异步流,每次GetAsyncEnumerator都会重新执行一次方法体,相当于把数据库查询再触发一遍;对于 Channel 的读取端,第一遍循环已经把数据消费完了,第二遍循环拿到的就是空集合。我遇到过一次线上问题,就是同事把同一个异步流传给了两个消费函数,第二个函数一直收到空数据,排查半天才发现数据在上一个循环里已经被“抽干”了。
解决思路很简单:想清楚这个流是单次的还是可重复的。单次流绝不能被多个消费者分享;如果确实需要遍历两遍,要么先ToListAsync()进内存,要么把数据源封装成工厂方法Func<IAsyncEnumerable<T>>,让每个消费者拿到的都是新的流。
5.2 DbContext 的生存期必须覆盖完整枚举
另一个典型错误是把 DbContext 提前释放:
public IAsyncEnumerable<Order> GetOrders() { using var db = new OrderDbContext(); // 危险 return db.Orders.AsAsyncEnumerable(); }调用方拿到异步流时,方法已经结束,DbContext 早就 dispose 了。等消费者真正开始MoveNextAsync(),就会得到一个 ObjectDisposedException。正确做法是把await foreach放进await using的代码块里,让 DbContext 的生存期和枚举的生存期对齐:
await using var db = new OrderDbContext(); await foreach (var order in db.Orders.AsAsyncEnumerable().WithCancellation(ct)) { await ProcessOrder(order, ct); }这个坑之所以隐蔽,是因为它不像普通方法那样在“调用时”就爆,而是在“第一次取数据时”才爆,出错点和赋值点离得远,定位起来很考验对惰性求值的理解。
5.3 普通 LINQ 不适用于 IAsyncEnumerable
标准System.Linq里的Where、Select是给IEnumerable<T>用的,直接套在异步流上会编译报错。初次接触的人十有八九会写:
var filtered = GetOrders().Where(o => o.Amount > 100); // 编译不过这不是语法问题,而是缺少对应的扩展方法。需要引入System.Linq.Async这个 NuGet 包:
using System.Linq.Async; await foreach (var order in GetOrders().WhereAwait(async o => await o.IsValidAsync())) { // ... }包里还有SelectAwait、ToArrayAsync、FirstAsync等一系列异步语义的 LINQ 操作。用到哪再引哪,别图省事把集合全都ToListAsync(),否则内存优势就丢光了。
5.4 阻塞等待是线程池杀手
有些人图省事,会把异步流塞进同步方法里用.Result或.GetAwaiter().GetResult()等待。我在代码评审里见过不止一次:
var first = GetOrders().FirstAsync().AsTask().GetAwaiter().GetResult(); // 反例在 ASP.NET Core 里,这样的阻塞等待会让请求线程在等待期间空转,线程池在高并发下可能被活活饿死。IAsyncEnumerable 的价值就在于整个消费链路都应该是异步的,所以从一开始就不要用同步方式桥接。
6. 我现在写异步流代码时的几条纪律
踩过一轮坑之后,我给自己定了几条很朴素的使用纪律。
第一,先判断数据是一次性流还是可重复流。凡是走数据库、网络、Channel 的数据流,一律按一次性流对待;需要重复消费的数据,显式地ToListAsync()一次并注明“这里就是要在内存里缓冲”。第二,资源持有者的生存期必须和枚举的生存期绑定。DbContext、HttpClient、Channel,都不允许在方法结束前被提前释放。第三,取消令牌永远不要省。迭代器内部要接收令牌,WithCancellation也要挂上,确保服务停机时能快速中断长循环。第四,跨库或者类库场景下,把ConfigureAwait(false)加上;有 UI 交互的场景,自己决定要不要保持同步上下文。第五,能流式就流式,要缓冲就明确缓冲,绝不偷偷把整张表搬进内存。
最后分享一个小经验。去年重构那个批处理服务时,我把原来ToListAsync()的那段改成了AsAsyncEnumerable()加 Channel 管道,没有引入任何重量级框架,只靠 BCL 自带的东西就把内存峰值压下来了。内存曲线从原来的锯齿状变成一条平缓的线,那一刻我才真正理解了“异步迭代”这四个字的含金量。如果你最近也在为大数据量处理发愁,我建议你先别急着上框架,试着把最痛的那条链路改成 IAsyncEnumerable,多半会有意外收获。