You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 06:05:19