【问题标题】:An enumerator wrapper that pre-buffers a number of items from underlying enumerator in advance一个枚举器包装器,它预先缓冲来自底层枚举器的多个项目
【发布时间】:2015-06-09 10:13:50
【问题描述】:

假设我有一些 IEnumerator<T>MoveNext() 方法中进行大量处理。

从该枚举器消耗的代码不仅消耗与可用数据一样快,而且偶尔会等待(其细节与我的问题无关)以同步需要恢复消耗的时间。但是当它下次调用MoveNext() 时,它需要尽可能快的数据。

一种方法是将整个流预先消耗到某个列表或数组结构中以进行即时枚举。然而,这会浪费内存,因为在任何一个时间点,只有一个项目正在使用,并且在整个数据无法放入内存的情况下,这将是令人望而却步的。

所以.net 中是否有一些通用的东西以某种方式包装一个枚举器/可枚举器,它预先异步预迭代底层枚举器几个项目并缓冲结果,以便它始终它的缓冲区中有许多可用的项目,并且调用 MoveNext 将永远不必等待?显然,消耗的项目,即由调用者的后续 MoveNext 迭代,将从缓冲区中删除。

注意我正在尝试做的部分事情也称为 Backpressure,并且,在 Rx 世界中,已经在 RxJava 中实现并且正在在Rx.NET 中讨论。 Rx(推送数据的 observables)可以被认为是枚举器的相反方法(枚举器允许拉取数据)。背压在拉取方法中相对容易,正如我的回答所示:只需暂停消费。推送时更难,需要额外的反馈机制。

【问题讨论】:

  • 这看起来像您之前的问题的副本,您说您会编辑:stackoverflow.com/questions/30700154/a-pre-buffering-enumerator
  • @Asad 我删除了旧问题并创建了这个问题,因为现有问题的 cmets 与当前形式的问题根本不匹配。我希望这个能让我更清楚我真正想要什么。
  • 这个问题似乎没有明显改变。您仍然存在这样的问题:一旦消费者耗尽缓冲区,在调用 MoveNext 时,预缓冲几个项目只会导致延迟两倍的频率。
  • 消费者通常不会用尽缓冲区,这实际上 在问题的第一个版本中缺少的重要限定。也就是说:The code consuming from that enumerator does not just consume as fast as data is available, but occasionally waits [..] But when it does the next call to MoveNext(), it needs the data as fast as possible.
  • @Asad 回复您的最后一条评论:您(和我的解决方案)有一​​个固定的缓冲区大小,这实际上足以解决我当前的问题,但不必修复它。预缓冲代码可以例如通过增加缓冲区的大小来适应消费者在长时间快速消费时威胁要完全耗尽缓冲区的消费者。使用 BlockingCollection(具有固定界限)很难做到的事情,但对于更动态的用例可能很有用,并且消除了消费者承诺固定大小的需要。

标签: c# .net enumerator backpressure


【解决方案1】:

您的自定义可枚举类的更简洁的替代方法是这样做:

public static IEnumerable<T> Buffer<T>(this IEnumerable<T> source, int bufferSize)
{
    var queue = new BlockingCollection<T>(bufferSize);

    Task.Run(() => {
        foreach(var i in source) queue.Add(i);
        queue.CompleteAdding();
    });

    return queue.GetConsumingEnumerable();
}

这可以用作:

var slowEnumerable = GetMySlowEnumerable();
var buffered = slowEnumerable.Buffer(10); // Populates up to 10 items on a background thread

【讨论】:

  • 很好的解决方案。我知道并使用过 BlockingCollection(并停止在需要快速的高并发代码中使用它,因为it is slow compared to ConcurrentQueue+AutoResetEvent),但我不知道GetConsumingEnumerable,在这个简单的生产者/消费者场景中,它看起来像是完美的解决方案.所以我毕竟不必自己动手。
  • @EugeneBeresovsky :) 。只是我觉得这里的最佳缓冲区大小是您最终计划消耗的任何内容(即缓冲所有内容,而不是前面的一小段距离)。我的意思是,如果你无论如何都要生成另一个线程,让它继续工作并尽可能提前做好准备。
  • @EugeneBeresovsky 不知道 BlockingCollection 很慢,但您可能可以使用ConcurrentQueue,并使用Add 锁定同步大小检查。
  • 不知道为什么您认为消费者无法确定合适的缓冲区大小。我的应用程序中有多达 10 个这样的枚举器,每个枚举器都会创建多达 10GB 的数据。这是以一定速度重播的数据(并且可以由用户暂停等)。计算必要的缓冲区大小取决于重播速度和项目/秒,因此很简单。您看不到的好处是 a) 极大地减少了内存消耗(与同时加载所有内容相比)和 b) 减少了延迟(在没有预缓冲的情况下按需读取时)。
  • @EugeneBeresovsky 我明白了。我想我只是把我自己的经验投射到这个上面。我从来没有使用过大到足以施加内存限制的数据集,但我经常不得不处理吞吐量问题。
【解决方案2】:

有不同的方法可以自己实现,我决定使用

  • 每个枚举器有一个专用线程,用于执行异步预缓冲
  • 要预缓冲的固定数量的元素

这对我手头的情况来说是完美的(只有少数,非常长时间运行的枚举器),但是例如如果您使用大量的枚举器,创建线程可能会太繁重,如果您需要更多动态的元素(可能基于项目的实际内容),固定数量的元素可能太不灵活。

到目前为止,我只测试了它的主要功能,可能还存在一些粗糙的边缘。可以这样使用:

int bufferSize = 5;
IEnumerable<int> en = ...;
foreach (var item in new PreBufferingEnumerable<int>(en, bufferSize))
{
    ...

这是枚举器的要点:

class PreBufferingEnumerator<TItem> : IEnumerator<TItem>
{
    private readonly IEnumerator<TItem> _underlying;
    private readonly int _bufferSize;
    private readonly Queue<TItem> _buffer;
    private bool _done;
    private bool _disposed;

    public PreBufferingEnumerator(IEnumerator<TItem> underlying, int bufferSize)
    {
        _underlying = underlying;
        _bufferSize = bufferSize;
        _buffer = new Queue<TItem>();
        Thread preBufferingThread = new Thread(PreBufferer) { Name = "PreBufferingEnumerator.PreBufferer", IsBackground = true };
        preBufferingThread.Start();
    }

    private void PreBufferer()
    {
        while (true)
        {
            lock (_buffer)
            {
                while (_buffer.Count == _bufferSize && !_disposed)
                    Monitor.Wait(_buffer);
                if (_disposed)
                    return;
            }
            if (!_underlying.MoveNext())
            {
                lock (_buffer)
                    _done = true;
                return;
            }
            var current = _underlying.Current; // do outside lock, in case underlying enumerator does something inside get_Current()
            lock (_buffer)
            {
                _buffer.Enqueue(current);
                Monitor.Pulse(_buffer);
            }
        }
    }

    public bool MoveNext()
    {
        lock (_buffer)
        {
            while (_buffer.Count == 0 && !_done && !_disposed)
                Monitor.Wait(_buffer);
            if (_buffer.Count > 0)
            {
                Current = _buffer.Dequeue();
                Monitor.Pulse(_buffer); // so PreBufferer thread can fetch more
                return true;
            }
            return false; // _done || _disposed
        }
    }

    public TItem Current { get; private set; }

    public void Dispose()
    {
        lock (_buffer)
        {
            if (_disposed)
                return;
            _disposed = true;
            _buffer.Clear();
            Current = default(TItem);
            Monitor.PulseAll(_buffer);
        }
    }

【讨论】:

    猜你喜欢
    • 2011-04-03
    • 1970-01-01
    • 2010-12-22
    • 2014-04-03
    • 1970-01-01
    • 2012-07-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多