【问题标题】:Reactive Extensions SelectMany with large objects带有大对象的反应式扩展 SelectMany
【发布时间】:2017-04-09 23:36:02
【问题描述】:

我有一小段代码可以模拟使用大型对象(即巨大的byte[])的流程。对于序列中的每一项,都会调用一个异步方法来获得一些结果。问题?事实上,它会抛出OutOfMemoryException。

与 LINQPad(C# 程序)兼容的代码:

void Main()
{
    var selectMany = Enumerable.Range(1, 100)
                   .Select(i => new LargeObject(i))
                   .ToObservable()
                   .SelectMany(o => Observable.FromAsync(() => DoSomethingAsync(o)));

    selectMany
        .Subscribe(r => Console.WriteLine(r));
}


private static async Task<int> DoSomethingAsync(LargeObject lo)
{
    await Task.Delay(10000);
    return lo.Id;
}

internal class LargeObject
{
    public int Id { get; }

    public LargeObject(int id)
    {
        this.Id = id;
    }

    public byte[] Data { get; } = new byte[10000000];
}

似乎它同时创建了所有对象。我怎样才能以正确的方式做到这一点?

基本思想是调用 DoSomethingAsync 以便为每个对象获取一些结果,这就是我使用 SelectMany 的原因。为了简化,我只是引入了一个Task.Delay,但在现实生活中它是一个可以同时处理一些项目的服务,所以我想引入一些并发机制来利用它。

请注意,理论上,一次处理少量项目不应该填满内存。实际上,我们只需要每个“大对象”来获取 DoSomethingAsync 方法的结果。在那之后,不再使用大对象。

【问题讨论】:

  • 我不知道你的问题是你的测试代码(Enumerable.Range 急切地创建所有大对象),还是你在生产中看到这个?无论哪种方式,如果某个序列创建了许多 LargeObjects 并且它们仍在使用中,所以不能被 GC 处理,那么是的,你会得到一个 OOM 异常。

标签: c# .net system.reactive reactive-programming


【解决方案1】:

我觉得我是repeating myself。与您上一个问题和我上一个回答类似,您需要做的是限制要同时创建的 bigObjects™ 的数量。

为此,您需要将对象创建和处理结合起来,并将其放在同一个线程池中。现在的问题是,我们使用异步方法来允许线程在我们的异步方法运行时做其他事情。由于您的慢速网络调用是异步的,因此您的(快速)对象创建代码将继续太快地创建大型对象。

相反,我们可以通过将对象创建与异步调用结合起来,使用 Rx 来记录并发运行的 Observable 的数量,并使用 .Merge(maxConcurrent) 来限制并发。

作为奖励,我们还可以设置查询执行的最短时间。只需 Zip 即可,延迟时间最短。

static void Main()
{
    var selectMany = Enumerable.Range(1, 100)
                        .ToObservable()
                        .Select(i => Observable.Defer(() => Observable.Return(new LargeObject(i)))
                            .SelectMany(o => Observable.FromAsync(() => DoSomethingAsync(o)))
                            .Zip(Observable.Timer(TimeSpan.FromMilliseconds(400)), (el, _) => el)
                        ).Merge(4);

    selectMany
        .Subscribe(r => Console.WriteLine(r));

    Console.ReadLine();
}


private static async Task<int> DoSomethingAsync(LargeObject lo)
{
    await Task.Delay(10000);
    return lo.Id;
}

internal class LargeObject
{
    public int Id { get; }

    public LargeObject(int id)
    {
        this.Id = id;
        Console.WriteLine(id + "!");
    }

    public byte[] Data { get; } = new byte[10000000];
}

【讨论】:

    【解决方案2】:

    好像是同时创建了所有的对象。

    是的,因为您是同时创建它们的。

    如果我简化你的代码,我可以告诉你原因:

    void Main()
    {
        var selectMany =
            Enumerable
                .Range(1, 5)
                .Do(x => Console.WriteLine($"{x}!"))
                .ToObservable()
                .SelectMany(i => Observable.FromAsync(() => DoSomethingAsync(i)));
    
        selectMany
            .Subscribe(r => Console.WriteLine(r));
    }
    
    private static async Task<int> DoSomethingAsync(int i)
    {
        await Task.Delay(1);
        return i;
    }
    

    运行它会产生:

    1! 2! 3! 4! 5! 4 3 5 2 1

    由于Observable.FromAsync,您允许源在任何结果返回之前运行完成。换句话说,您正在快速构建所有大型对象,但在缓慢地处理它们。

    您应该允许 Rx 同步运行,但在默认调度程序上运行,这样您的主线程就不会被阻塞。然后代码将在没有任何内存问题的情况下运行,并且您的程序将在主线程上保持响应。

    下面是代码:

    var selectMany =
        Observable
            .Range(1, 100, Scheduler.Default)
            .Select(i => new LargeObject(i))
            .Select(o => DoSomethingAsync(o))
            .Select(t => t.Result);
    

    (我已经用Observable.Range(1, 100) 有效地替换了Enumerable.Range(1, 100).ToObservable(),因为这也有助于解决一些问题。)

    我已尝试测试其他选项,但到目前为止,任何允许 DoSomethingAsync 异步运行的操作都会遇到内存不足错误。

    【讨论】:

    • 感谢您的回答,@Enigmativity,但我想我错过了一些东西。我正在调用的异步方法是一个可以同时处理项目的远程服务。在处理另一个项目之前等待一个项目被处理并不是最佳选择。您认为我可以同时处理多个项目(3 个或 4 个)以利用并发性而不会遇到内存问题吗?
    • @SuperJMN,如果您尝试创建的 LargeObjects 超出您的分配范围,您将收到 OOM 异常。这与 Rx 有关。如果它们需要独立于正在执行的DoSomethingAsync 创建,那么你就有麻烦了。在我看来,您实际上希望将这些保留到队列中并脱离 Rx。
    • @SuperJMN - 我尝试了很多方法来限制处理,但它只是没有用。一次处理一个对象的计算效率不高,但它具有内存效率。这取决于您要达到的效率。 Shlomo 的回答对于计算效率来说甚至更糟。如果我想到什么我会告诉你的。
    • 谢谢!在 Stack Overflow 之外,有人建议我使用 TPL/Workflows,但我不知道如何处理。
    • @LeeCampbell 你的意思是在 DoSomethingAsync 方法中加载/卸载大对象吗?关于队列,我看不到在哪里使用。你能发布一些例子吗?
    【解决方案3】:

    ConcatMap 开箱即用地支持这一点。我知道这个操作符在 .net 中不可用,但是你可以使用 Concat 操作符来做同样的事情,它将订阅每个内部源直到前一个完成。

    【讨论】:

    【解决方案4】:

    您可以通过这种方式引入时间间隔延迟:

    var source = Enumerable.Range(1, 100)
       .ToObservable()
       .Zip(Observable.Interval(TimeSpan.FromSeconds(1)), (i, ts) => i)
       .Select(i => new LargeObject(i))
       .SelectMany(o => Observable.FromAsync(() => DoSomethingAsync(o)));
    

    因此,不是一次提取所有 100 个整数,而是立即将它们转换为 LargeObject,然后在所有 100 个上调用 DoSomethingAsync,而是将整数逐一滴出,每个整数间隔一秒。


    这就是 TPL+Rx 解决方案的样子。不用说它不如单独的 Rx 或单独的 TPL 优雅。但是,我认为这个问题不太适合 Rx:

    void Main()
    {
        var source = Observable.Range(1, 100);
    
        const int MaxParallelism = 5;
        var transformBlock = new TransformBlock<int, int>(async i => await DoSomethingAsync(new LargeObject(i)),
            new ExecutionDataflowBlockOptions { MaxDegreeOfParallelism = MaxParallelism });
        source.Subscribe(transformBlock.AsObserver());
        var selectMany = transformBlock.AsObservable();
    
        selectMany
            .Subscribe(r => Console.WriteLine(r));
    }
    

    【讨论】:

    • 我可以理解这在实践中可能有效,但是选择一秒延迟是任意的,并且仍然可能出现内存不足错误,或者它会显着减慢计算速度。这不是一个可靠的解决方案。
    • 编辑添加 TPL 答案。 Rx 在这里不亮。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-02-10
    • 2011-08-27
    相关资源
    最近更新 更多