EF Core使用IAsyncEnumerable流式查询长时间运行抛出异常
我使用IAsyncEnumerable从数据库流式传输超大规模表的行到应用程序,处理后写入CSV文件流。由于表内行数极多,查询需运行数小时,频繁遇到如下异常:
System.InvalidOperationException Invalid attempt to call CheckDataIsReady when reader is closed.
System.InvalidOperationException: Invalid attempt to call CheckDataIsReady when reader is closed.
at Microsoft.Data.SqlClient.SqlDataReader.CheckDataIsReady(Int32 columnIndex, Boolean allowPartiallyReadColumn, Boolean permitAsync, String methodName)
at Microsoft.Data.SqlClient.SqlDataReader.GetFieldValueInternal[T](Int32 i)
at Microsoft.Data.SqlClient.SqlDataReader.GetFieldValue[T](Int32 i)
at lambda_method965(Closure , QueryContext , DbDataReader , ResultContext , SingleQueryResultCoordinator )
at Microsoft.EntityFrameworkCore.Query.Internal.SingleQueryingEnumerable1.Enumerator.MoveNext() at System.Linq.AsyncEnumerable.AsyncEnumerableAdapter1.MoveNextCore() in //Ix.NET/Source/System.Linq.Async/System/Linq/Operators/ToAsyncEnumerable.cs:line 79
at System.Linq.AsyncIteratorBase1.MoveNextAsync() in /_/Ix.NET/Source/System.Linq.Async/System/Linq/AsyncIterator.cs:line 77 at System.Linq.AsyncIteratorBase1.MoveNextAsync() in //Ix.NET/Source/System.Linq.Async/System/Linq/AsyncIterator.cs:line 77
我的代码如下:
await context.Database .CreateExecutionStrategy() .ExecuteInTransactionAsync(async cancellationToken => { var entityResult = context.Set<TEntity>().AsNoTracking().ToAsyncEnumerable(); var done = false; await using var enumerator = entityResult.GetAsyncEnumerator(); await using var stream = new MemoryStream(); await using var writer = new StreamWriter(stream); var csv = new CsvWriter(writer, CultureInfo.InvariantCulture, true); csv.Context.RegisterClassMap(new EntityClassMap<TEntity>()); csv.WriteHeader<TEntity>(); csv.NextRecord(); while (await enumerator.MoveNextAsync()) // Cannot use foreach, because of some other stuff below { var entity = enumerator.Current; csv.WriteRecord(entity); csv.NextRecord(); // some other stuff } }, _ => Task.FromResult(true), // We are just reading, so we can always commit the transaction System.Data.IsolationLevel.ReadUncommitted, // Do not block the whole table while reading. This is essentially the same as WITH(NOLOCK). cancellationToken);
推测问题源于网络波动或数据库繁忙,需要弹性处理机制,但无法使用默认的SqlServerRetryingExecutionStrategy,因其会将所有行缓冲至内存(达数百GB),寻求解决方案。
1. 分段读取+断点续传
将全表查询拆分为多个小批次,通过排序键(自增ID、时间戳等)分段读取,同时记录已处理的最后位置,异常时从断点恢复,避免一次性加载全表:
- 选择连续的排序字段作为分段依据,每次查询仅获取
排序键 > 上次最后值的批次数据 - 持久化断点记录(写入本地文件或小型数据库表),每次成功处理完一个批次后更新断点
- 示例代码片段:
long lastProcessedId = LoadLastProcessedId(); // 从本地文件/数据库加载断点 bool hasMore = true; const int batchSize = 1000; while (hasMore && !cancellationToken.IsCancellationRequested) { try { var batch = await context.Set<TEntity>() .AsNoTracking() .Where(e => e.Id > lastProcessedId) .OrderBy(e => e.Id) .Take(batchSize) .ToAsyncEnumerable() .ToListAsync(cancellationToken); if (batch.Count == 0) { hasMore = false; break; } foreach (var entity in batch) { csv.WriteRecord(entity); csv.NextRecord(); // 处理其他业务逻辑 lastProcessedId = entity.Id; } await writer.FlushAsync(cancellationToken); SaveLastProcessedId(lastProcessedId); // 持久化最新断点 } catch (SqlException ex) when (IsTransientError(ex)) // 仅重试临时错误 { await Task.Delay(TimeSpan.FromSeconds(2), cancellationToken); continue; } }
2. 自定义轻量批次重试策略
针对单个批次的查询和处理实现重试逻辑,仅在临时错误(网络超时、数据库繁忙)时重试,避免全表缓冲:
private async Task<T> RetryOnTransientErrorsAsync<T>(Func<Task<T>> operation, int maxRetries = 3) { int retryCount = 0; while (true) { try { return await operation(); } catch (SqlException ex) when (IsTransientError(ex) && retryCount < maxRetries) { retryCount++; // 指数退避重试 await Task.Delay(TimeSpan.FromSeconds(Math.Pow(2, retryCount)), cancellationToken); } } } // 判断是否为SQL临时错误 private bool IsTransientError(SqlException ex) { return ex.Errors.Cast<SqlError>().Any(error => error.Number is 4060 or 10928 or 10929 or 40197 or 40501 or 40613 or 49918 or 49919 or 49920); }
3. 优化数据库连接与查询参数
调整连接字符串和EF Core配置,减少连接异常概率:
- 在连接字符串中添加
ConnectRetryCount=5;ConnectRetryInterval=2,启用连接层面的重试 - 设置
CommandTimeout=3600(或更大值),避免单个批次查询超时 - 去掉外层事务包裹,直接用
WITH(NOLOCK)查询提示替代ReadUncommitted事务,长时间事务易引发连接问题:
var entityResult = context.Set<TEntity>() .AsNoTracking() .FromSqlRaw("SELECT * FROM YourTable WITH(NOLOCK)") .ToAsyncEnumerable();
内容的提问来源于stack exchange,提问作者wertzui

