【问题标题】:Async with huge data streams与大量数据流异步
【发布时间】:2014-09-17 21:58:25
【问题描述】:

我们使用 IEnumerables 从数据库中返回大量数据集:

public IEnumerable<Data> Read(...)
{
    using(var connection = new SqlConnection(...))
    {
        // ...
        while(reader.Read())
        {
            // ...
            yield return item;
        }
    }
}

现在我们想使用异步方法来做同样的事情。但是异步没有 IEnumerables,所以我们必须将数据收集到一个列表中,直到加载整个数据集:

public async Task<List<Data>> ReadAsync(...)
{
    var result = new List<Data>();
    using(var connection = new SqlConnection(...))
    {
        // ...
        while(await reader.ReadAsync().ConfigureAwait(false))
        {
            // ...
            result.Add(item);
        }
    }
    return result;
}

这会消耗服务器上的大量资源,因为所有数据必须在列表中才能返回。 IEnumerables 处理大型数据流的最佳且易于使用的异步替代方案是什么?我想避免在处理时将所有数据存储在内存中。

【问题讨论】:

  • 这会消耗服务器上的大量资源... 那么在这种情况下,服务器执行环境是什么? (例如 ASP.NET、WCF 服务等)什么是客户端执行环境? (网络浏览器、富客户端 .NET 应用等)
  • Reactive Extensions 包括 Async Enumerable's,你应该会发现它很有帮助。
  • @ChristopherHarris,你的意思是来自 Ineteractive Extensions (Ix) 的 IAsyncEnumerable 吗?可用的信息很少,例如these slidesthis blog。除了 Ix Experimental,还有其他版本吗?

标签: c# .net task-parallel-library async-await


【解决方案1】:

最简单的选择是使用TPL Dataflow。您需要做的就是配置一个ActionBlock 来处理处理(如果您愿意,可以并行处理)并异步“发送”项目一个一个地发送到其中。
我还建议设置一个BoundedCapacity,当处理无法处理速度时,它将限制读取器从数据库中读取。

var block = new ActionBlock<Data>(
    data => ProcessDataAsync(data),
    new ExecutionDataflowBlockOptions
    {
        BoundedCapacity = 1000,
        MaxDegreeOfParallelism = Environment.ProcessorCount
    });

using(var connection = new SqlConnection(...))
{
    // ...
    while(await reader.ReadAsync().ConfigureAwait(false))
    {
        // ...
       await block.SendAsync(item);
    }
}

您也可以使用Reactive Extensions,但这是一个比您可能需要的更复杂和更强大的框架。

【讨论】:

  • 很好的回答其他解决方案,包括我自己编写的解决方案,无缘无故地复杂得多。 +1
  • 数据流的胜利。它的多功能性真的令人惊讶。我刚刚编写了一个无线电服务器,它真正受益于它提供的缓冲区边界。爱它。 +1
  • @l3arnon,您能否澄清一下它如何帮助解决这个问题:这将消耗服务器上的大量资源,因为所有数据在返回之前都必须在列表中。无论您是否在服务器上使用 TPL Dataflow 来限制异步数据读取,当所有数据都已获取时,响应都会发送到客户端。除非您使用 SignalR 或 WCF 流式传输。
  • @Noseratio OP 没有提到客户,也没有在代码中解决这个问题。她/他询问有关数据库检索和处理的问题。客户端服务器架构的解决方案取决于客户端使用的技术和所需的行为。
  • item 在哪里,请添加更多代码以及BoundedCapacity = 1000 如何在这里发挥作用.. 努力追随,ProcessDataAsync 是什么,它的类型有 SendAsyncSendAsync ActionBlock 是如何被调用的。,你能不能添加 cmets 并使其更容易理解,这很令人困惑。对于那些不熟悉TPL Dataflow的人
【解决方案2】:

大多数时候在处理 async/await 方法时,我发现更容易解决问题,并使用函数 (Func&lt;...&gt;) 或操作 (Action&lt;...&gt;) 而不是临时代码,尤其是使用 @987654323 @ 和yield

换句话说,当我想到“异步”时,我试图忘记函数“返回值”的旧概念,否则它是如此明显且我们如此熟悉。

例如,如果您将初始同步代码更改为此(processor 是最终将执行您对一个数据项所做的代码):

public void Read(..., Action<Data> processor)
{
    using(var connection = new SqlConnection(...))
    {
        // ...
        while(reader.Read())
        {
            // ...
            processor(item);
        }
    }
}

那么,异步版本写起来就很简单了:

public async Task ReadAsync(..., Action<Data> processor)
{
    using(var connection = new SqlConnection(...))
    {
        // note you can use connection.OpenAsync()
        // and command.ExecuteReaderAsync() here
        while(await reader.ReadAsync())
        {
            // ...
            processor(item);
        }
    }
}

如果您可以通过这种方式更改代码,则不需要任何扩展或额外的库或 IAsyncEnumerable 的东西。

【讨论】:

  • 我可能会为 asyncProcessor 添加一个选项:public async void ReadAsync(..., Func&lt;Data,Task&gt; asyncProcessor) 并使用它:await asyncProcessor(item);
  • +1,不确定我是否会使用这个解决方案,但它看起来像是跳出框框思考,并为我提供了解决其他问题的想法
【解决方案3】:

这将消耗服务器上的大量资源,因为所有 返回前数据必须在列表中。什么是最好的和容易的 使用 IEnumerables 的异步替代方法来处理大数据 流?我想避免将所有数据存储在内存中,同时 处理。

如果您不想一次将所有数据发送到客户端,您可以考虑使用Reactive Extensions (Rx)(在客户端)和SignalR(在客户端和服务器上)来处理。

SignalR 允许异步向客户端发送数据。 Rx 将允许在数据项到达客户端时将 LINQ 应用于异步数据项序列。但是,这会改变您的客户端-服务器应用程序的整个代码模型。

示例(Samuel Jack 的博客文章):

相关问题(如果不重复):

【讨论】:

  • 我当然不会反对(我认为链接的答案很棒,因为它是唯一不引入外部依赖的解决方案),但您似乎回答了 更多 个问题而不是实际上是被问到的。我可以看到“服务器”这个词可能会让您绊倒,但最终整个事情似乎归结为“我如何实现异步yield return”(您的链接答案完全解决了这个问题)而不是“如何实现完整的客户端-服务器流”。
  • @KirillShlenskiy,无论客户端/服务器方面如何,异步 yield return 正是 Rx 完美解决的问题,IMO。 IObservable&lt;T&gt; 在几乎每一个方面都与 IEnumerable&lt;T&gt; 有着深刻的双重性。
  • ...实际上,再想一想,在当前问题的上下文中,有一个关键方面可能很重要,而 Rx 没有解决(由于缺少 OnNextAsync 和消费者 - > 生产者反馈机制):背压。在消费者比生产者慢的情况下,拥有一个容量有限的切换点对于防止中间缓存中的项目堆积非常宝贵(“消耗大量资源”正如 user1224129 所说),这将使 TPL Dataflow 或您自己的链接解决方案更适合。
  • @KirillShlenskiy,Rx 也为此提供了缓冲 (Observable.Buffer),您也可以跳过多余的项目。要点是,使用 Dataflow,您必须先构建一个有限序列,然后才能使用 LINQ 查询对其进行处理。使用 Rx,您不需要。相关:stackoverflow.com/a/24245474/1768303
  • @KirillShlenskiy 如果她/他有兴趣这样做,OP 可以轻松地将大多数(如果不是全部)LINQ 查询与等效的 Dataflow 块交换。
【解决方案4】:

正如其他一些海报所提到的,这可以通过 Rx 来实现。使用 Rx,该函数将返回一个可以订阅的 IObservable&lt;Data&gt;,并在数据可用时将数据推送给订阅者。 IObservable 也支持 LINQ 并添加了自己的一些扩展方法。

更新

我添加了几个通用帮助方法,以使阅读器的使用可重用并支持取消。

public static class ObservableEx
    {
        public static IObservable<T> CreateFromSqlCommand<T>(string connectionString, string command, Func<SqlDataReader, Task<T>> readDataFunc)
        {
            return CreateFromSqlCommand(connectionString, command, readDataFunc, CancellationToken.None);
        }

        public static IObservable<T> CreateFromSqlCommand<T>(string connectionString, string command, Func<SqlDataReader, Task<T>> readDataFunc, CancellationToken cancellationToken)
        {
            return Observable.Create<T>(
                async o =>
                {
                    SqlDataReader reader = null;

                    try
                    {                        
                        using (var conn = new SqlConnection(connectionString))
                        using (var cmd = new SqlCommand(command, conn))
                        {
                            await conn.OpenAsync(cancellationToken);
                            reader = await cmd.ExecuteReaderAsync(CommandBehavior.CloseConnection, cancellationToken);

                            while (await reader.ReadAsync(cancellationToken))
                            {
                                var data = await readDataFunc(reader);
                                o.OnNext(data);                                
                            }

                            o.OnCompleted();
                        }
                    }
                    catch (Exception ex)
                    {
                        o.OnError(ex);
                    }

                    return reader;
                });
        }
    }

ReadData 的实现现在大大简化了。

     private static IObservable<Data> ReadData()
    {
        return ObservableEx.CreateFromSqlCommand(connectionString, "select * from Data", async r =>
        {
            return await Task.FromResult(new Data()); // sample code to read from reader.
        });
    }

用法

您可以通过给 Observable 一个 IObserver 来订阅它,但也有需要 lambda 的重载。随着数据可用,将调用 OnNext 回调。如果出现异常,将调用 OnError 回调。最后,如果没有更多数据,则调用 OnCompleted 回调。

如果您想取消 observable,只需处置订阅即可。

void Main()
{
   // This is an asyncrhonous call, it returns straight away
    var subscription = ReadData()
        .Skip(5)                        // Skip first 5 entries, supports LINQ               
        .Delay(TimeSpan.FromSeconds(1)) // Rx operator to delay sequence 1 second
        .Subscribe(x =>
    {
        // Callback when a new Data is read
        // do something with x of type Data
    },
    e =>
    {
        // Optional callback for when an error occurs
    },
    () =>
    {
        //Optional callback for when the sequenc is complete
    }
    );

    // Dispose subscription when finished
    subscription.Dispose();

    Console.ReadKey();
}

【讨论】:

    【解决方案5】:

    我认为 Rx 绝对是在这种情况下要走的路,因为可观察序列是可枚举序列的形式对偶。

    如上一个答案中所述,您可以从头开始将序列重写为可观察对象,但也有几种方法可以继续编写迭代器块,然后异步展开它们。

    1) 只需将 enumerable 转换为 observable,如下所示:

    using System.Reactive.Linq;
    using System.Reactive.Concurrency;
    
    var enumerable = Enumerable.Range(10);
    var observable = enumerable.ToObservable();
    var subscription = observable.Subscribe(x => Console.WriteLine(x));
    

    这将使您的可枚举通过将其通知推送到任何下游观察者来表现得像一个可观察的。在这种情况下,当调用 Subscribe 时,它​​会同步阻塞,直到处理完所有数据。如果您希望它完全异步,可以使用以下命令将其设置为不同的线程:

    var observable = enumerable.ToObservable().SubscribeOn(NewThreadScheduler.Default);
    

    现在可枚举的展开将在一个新线程中完成,订阅方法将立即返回。

    2) 使用另一个异步事件源展开可枚举:

    var enumerable = Enumerable.Range(10);
    var observable = Observable.Timer(TimeSpan.Zero, TimeSpan.FromSeconds(1))
                               .Zip(enumerable, (t, x) => x);
    var subscription = observable.Subscribe(x => Console.WriteLine(x));
    

    在这种情况下,我设置了一个计时器,每秒触发一次,每当它触发时,它就会向前移动迭代器。现在计时器可以很容易地被任何事件源替换,以准确控制迭代器何时向前移动。

    我发现自己喜欢迭代器块的语法和语义(例如 try/finally 块和 dispose 会发生什么),因此即使在设计异步操作时我也会偶尔使用这些设计。

    【讨论】:

      猜你喜欢
      • 2022-01-03
      • 2013-01-17
      • 1970-01-01
      • 2021-08-06
      • 2013-11-22
      • 1970-01-01
      • 1970-01-01
      • 2020-04-30
      • 1970-01-01
      相关资源
      最近更新 更多