
接手维护一个老批处理服务之后我对 IAsyncEnumerable 的态度从“知道有这东西”变成了“离不了它”。那个服务每天凌晨会从数据库拉一张几百万行的订单表做统计原来的实现是TaskListOrder的经典套路一个ToListAsync()把所有数据先装进内存再说。数据量小时一切正常到了线上真实数据量进程 GC 频率高得离谱CPU 大量花在清理托管堆上。后来把核心读取改成异步流内存峰值降了一个数量级那份快乐我现在都记得。这篇文章我不打算把 API 文档念一遍而是按我自己理解它的顺序讲先说你到底在解决什么问题再看接口和编译器的底层设计然后给几个能直接抄进生产代码的用法最后说几个特别容易踩的坑。1. 异步迭代出现的真正理由别再把整张表“搬”进内存1.1 传统写法的三个隐患早些年写数据批处理我们习惯性先拿一个集合ListOrder orders await _db.Orders .Where(o o.Status Paid) .ToListAsync();这句话背后藏着三个问题。第一数据库把所有匹配记录一次性全部返回网络传输、数据读取都是从头等到尾。第二EF Core 把每一行打成实体对象再塞进ListOrder这个 List 在整个方法生命周期甚至更久都躺在内存里。第三后续任何一次遍历、聚合、筛选都要先跨过这块已经占好的“大石头”。我在生产里见过最夸张的例子是某个统计服务把 800 万条明细ToList()之后再聚合内存 GC 从几秒一次变成几百毫秒一次最后整个服务 CPU 卡在垃圾回收上业务吞吐量直接掉到原来的三分之一。问题不在业务逻辑而在“数据还没用上就已经全部常驻内存”。1.2 TaskIEnumerable 只是“半截”异步很多人觉得TaskIEnumerableT已经是异步了但它的本质只是“等一下再拿集合”。拿到集合之后你对它做foreach仍然是同步遍历如果这个集合是惰性的每一次MoveNext()都可能阻塞调用线程。也就是说第一段异步只解决了“到达数据前不占线程”的问题没有解决“拿到数据流之后逐条读取不阻塞”的问题。IEnumerableT本身能流式读取但它读取的过程是同步的。当数据源真的是网络接口、远程 API、数据库游标这类需要异步 I/O 的地方你要么让线程傻等要么手动拆成一页一页去拉。而 IAsyncEnumerable 想解决的问题就是“异步地从远处拿数据拿到一条就处理一条拿不到的时候线程去干别的”。1.3 三种数据形态放在一起看形态数据产生方式内存表现等待时是否占用线程IEnumerableT同步拉取可以惰性但读取时阻塞占用调用线程TaskIEnumerableT异步拿到集合遍历仍同步通常一次性全量入内存等待期间不占遍历时占IAsyncEnumerableT逐个异步拉取边取边甩天然流式等待期间不占也不阻塞这个对比基本就是我在选型时的判断依据。数据量只有几十条时怎么写都无所谓一旦数据可能是几万、几百上千万条IAsyncEnumerable 的价值就立刻体现出来。2. 接口和编译器的三个关键设计Current、MoveNextAsync 与 ValueTask2.1 接口本身很短信息量不小IAsyncEnumerable 的接口定义非常简洁public interface IAsyncEnumerableout T { IAsyncEnumeratorT GetAsyncEnumerator(CancellationToken cancellationToken default); } public interface IAsyncEnumeratorout T : IAsyncDisposable { T Current { get; } ValueTaskbool MoveNextAsync(); }整个抽象就三件事Current给出当前值MoveNextAsync()异步判断还有没有下一个DisposeAsync()负责清理。GetAsyncEnumerator接收一个取消令牌让外部可以在枚举开始前就挂上取消信号。注意泛型参数带out意味着IAsyncEnumerableDog可以直接当成IAsyncEnumerableAnimal用这一点和IEnumerableT的协变语义一致定义接口时非常有用。2.2 yield 和 await 混在一起编译器怎么处理在 C# 8 里写异步迭代器最直观的形态长这样async IAsyncEnumerableint 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()的返回类型不是Taskbool而是ValueTaskbool这个细节值得单独说。异步流的特点是很多情况下下一个元素已经准备好了。比如 Channel 缓冲区里正好有数据或者包装的同步集合里还有剩余项这时候MoveNextAsync()完全可以同步返回true根本不需要一次真正的异步等待。如果接口设计成Taskbool每取一个元素都可能要分配一个 Task 对象。一百万个元素就有一百万次分配GC 压力相当大。ValueTask的设计思路是“能同步返回就同步返回需要真正异步时才去等待”在零分配和高吞吐之间做了很好的平衡。虽然日常使用中你不会感知到它但这个选择直接关系到百万行数据流的稳定性。3. await foreach 用法里的几个细节取消、配置与资源释放3.1 消费一个异步流的基本逻辑消费 IAsyncEnumerable 的常规入口是await foreachawait 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 IAsyncEnumerableOrder 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 IAsyncEnumerableOrder 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 流式读取AsAsyncEnumerableEF Core 3.0 以后IQueryableT上可以直接调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()直接返回IAsyncEnumerableTvar channel Channel.CreateUnboundedOrder(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可以直接吃 IAsyncEnumerableawait 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()进内存要么把数据源封装成工厂方法FuncIAsyncEnumerableT让每个消费者拿到的都是新的流。5.2 DbContext 的生存期必须覆盖完整枚举另一个典型错误是把 DbContext 提前释放public IAsyncEnumerableOrder 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是给IEnumerableT用的直接套在异步流上会编译报错。初次接触的人十有八九会写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多半会有意外收获。