【问题标题】:Multiple threads accesing IEnumerable using yield多个线程使用 yield 访问 IEnumerable
【发布时间】:2012-10-26 14:32:30
【问题描述】:

我正在使用第三方库来迭代一些非常大的平面文件,这可能需要几分钟。该库提供了一个枚举器,因此您可以在枚举器提取平面文件中的下一项时产生并处理每个结果。

例如:

IEnumerable<object> GetItems()
{
    var cursor = new Cursor;

    try
    {
        cursor.Open();

        while (!cursor.EOF)
        {
            yield return new //object;

            cursor.MoveNext();
        }

    }
    finally
    {
        if (cursor.IsOpen)
        {
            cursor.Close();
        }
    }
}

我想要实现的是拥有相同 Enumerable 的两个消费者,因此我不必两次提取信息,因此每个消费者仍然可以在每个项目到达时对其进行处理,而无需一直等待马上到达。

IEnumerable<object> items = GetItems();

new Thread(SaveToDateBase(items)).Start();
new Thread(SaveSomewhereElse(items)).Start();

我想我想要达到的目标是

“如果消费者要求的项目已经被提取,则放弃它,否则移动下一步并等待”但我意识到两个线程之间可能存在 MoveNext() 冲突。

如果没有任何关于如何实现的想法,这样的事情是否已经存在?

谢谢

【问题讨论】:

    标签: c# ienumerable ienumerator


    【解决方案1】:

    Pipelines pattern implementation 使用 .NET 4 BlockingCollection&lt;T&gt; 和 TPL 任务是您正在寻找的。用完整的例子查看我的答案in this StackOverflow post

    示例:3 个同时使用的消费者

    BlockingCollection<string> queue = new BlockingCollection<string>();    
    public void Start()
    {
        var producerWorker = Task.Factory.StartNew(() => ProducerImpl());
        var consumer1 = Task.Factory.StartNew(() => ConsumerImpl());
        var consumer2 = Task.Factory.StartNew(() => ConsumerImpl());
        var consumer3 = Task.Factory.StartNew(() => ConsumerImpl());
    
        Task.WaitAll(producerWorker, consumer1, consumer2, consumer3);
    }
    
    private void ProducerImpl()
    {
       // 1. Read a raw data from a file
       // 2. Preprocess it
       // 3. Add item to a queue
       queue.Add(item);
    }
    
    // ConsumerImpl must be thrad safe 
    // to allow launching multiple consumers simulteniously
    private void ConsumerImpl()
    {
        foreach (var item in queue.GetConsumingEnumerable())
        {
            // TODO
        }
    }
    

    如果还有不清楚的地方,请告诉我。

    管道流程的高级图:

    【讨论】:

    • 现在只是看看,但我无法将您的“TPL 特定管道实施”解决方案应用于我上面的问题。
    • @MarkVickery:您在寻找其他样本吗?还是已经找到解决方案?
    • 很快会进一步查看并返回更新和/或接受的答案。
    • 我查看了 BlockingCollection,虽然它确实解决了另一个问题,但我认为它并没有解决发布的问题。非常感谢。
    • 这里的问题是每件商品只由三个消费者之一处理。 OP 的示例是,有一个生产者,然后是两个消费者,他们各自处理每个项目,但并行进行消费,而不是一个管道,一个在另一个之前处理项目。这就是为什么您的解决方案不适用的原因。
    【解决方案2】:

    基本上你想要的是缓存IEnumerable&lt;T&gt; 的数据,但在存储之前无需等待它完成。你可以这样做:

    public static IEnumerable<T> Cache<T>(this IEnumerable<T> source)
    {
        return new CacheEnumerator<T>(source);
    }
    
    private class CacheEnumerator<T> : IEnumerable<T>
    {
        private CacheEntry<T> cacheEntry;
        public CacheEnumerator(IEnumerable<T> sequence)
        {
            cacheEntry = new CacheEntry<T>();
            cacheEntry.Sequence = sequence.GetEnumerator();
            cacheEntry.CachedValues = new List<T>();
        }
    
        public IEnumerator<T> GetEnumerator()
        {
            if (cacheEntry.FullyPopulated)
            {
                return cacheEntry.CachedValues.GetEnumerator();
            }
            else
            {
                return iterateSequence<T>(cacheEntry).GetEnumerator();
            }
        }
    
        IEnumerator IEnumerable.GetEnumerator()
        {
            return this.GetEnumerator();
        }
    }
    
    private static IEnumerable<T> iterateSequence<T>(CacheEntry<T> entry)
    {
        for (int i = 0; entry.ensureItemAt(i); i++)
        {
            yield return entry.CachedValues[i];
        }
    }
    
    private class CacheEntry<T>
    {
        public bool FullyPopulated { get; private set; }
        public IEnumerator<T> Sequence { get; set; }
    
        //storing it as object, but the underlying objects will be lists of various generic types.
        public List<T> CachedValues { get; set; }
    
        private static object key = new object();
        /// <summary>
        /// Ensure that the cache has an item a the provided index.  If not, take an item from the 
        /// input sequence and move to the cache.
        /// 
        /// The method is thread safe.
        /// </summary>
        /// <returns>True if the cache already had enough items or 
        /// an item was moved to the cache, 
        /// false if there were no more items in the sequence.</returns>
        public bool ensureItemAt(int index)
        {
            //if the cache already has the items we don't need to lock to know we 
            //can get it
            if (index < CachedValues.Count)
                return true;
            //if we're done there's no race conditions hwere either
            if (FullyPopulated)
                return false;
    
            lock (key)
            {
                //re-check the early-exit conditions in case they changed while we were
                //waiting on the lock.
    
                //we already have the cached item
                if (index < CachedValues.Count)
                    return true;
                //we don't have the cached item and there are no uncached items
                if (FullyPopulated)
                    return false;
    
                //we actually need to get the next item from the sequence.
                if (Sequence.MoveNext())
                {
                    CachedValues.Add(Sequence.Current);
                    return true;
                }
                else
                {
                    Sequence.Dispose();
                    FullyPopulated = true;
                    return false;
                }
            }
        }
    }
    

    示例用法:

    private static IEnumerable<int> interestingIntGenertionMethod(int maxValue)
    {
        for (int i = 0; i < maxValue; i++)
        {
            Thread.Sleep(1000);
            Console.WriteLine("actually generating value: {0}", i);
            yield return i;
        }
    }
    
    public static void Main(string[] args)
    {
        IEnumerable<int> sequence = interestingIntGenertionMethod(10)
            .Cache();
    
        int numThreads = 3;
        for (int i = 0; i < numThreads; i++)
        {
            int taskID = i;
            Task.Factory.StartNew(() =>
            {
                foreach (int value in sequence)
                {
                    Console.WriteLine("Task: {0} Value:{1}",
                        taskID, value);
                }
            });
        }
    
        Console.WriteLine("Press any key to exit...");
        Console.ReadKey(true);
    }
    

    【讨论】:

      猜你喜欢
      • 2011-01-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-12-03
      • 1970-01-01
      • 2016-11-07
      • 1970-01-01
      相关资源
      最近更新 更多