【问题标题】:Cannot retrieve chunks of data from BlockingCollection<T>无法从 BlockingCollection<T> 检索数据块
【发布时间】:2023-01-25 09:40:41
【问题描述】:

我经常发现自己确实想分块而不是一个接一个地传输数据。通常我在需要执行一些基于 I/O 的操作时这样做,比如我想限制往返的数据库插入。 所以我得到了这个不错的小扩展方法:

        public static IEnumerable<List<T>> Split<T>(this IEnumerable<T> data, int size)
        {            
            using (var enumerator = data.GetEnumerator())
            {
                while (enumerator.MoveNext())
                {
                    yield return YieldBatchElements(enumerator, size - 1).ToList();
                }
            }

            IEnumerable<TU> YieldBatchElements<TU>(
                IEnumerator<TU> source,
                int batchSize)
            {
                yield return source.Current;
                for (var i = 0; i < batchSize && source.MoveNext(); i++)
                {
                    yield return source.Current;
                }
            }
        }

这很好用,但我注意到它不适用于BlockCollection&lt;T&gt; GetConsumingEnumerable

我创建了以下测试方法来证明我的发现:

        [Test]
        public static void ConsumeTest()
        {
            var queue = new BlockingCollection<int>();
            var i = 0;
            foreach (var x in Enumerable.Range(0, 10).Split(3))
            {
                Console.WriteLine($"Fetched chunk: {x.Count}");
                Console.WriteLine($"Fetched total: {i += x.Count}");
            }
            //Fetched chunk: 3
            //Fetched total: 3
            //Fetched chunk: 3
            //Fetched total: 6
            //Fetched chunk: 3
            //Fetched total: 9
            //Fetched chunk: 1
            //Fetched total: 10
         

            Task.Run(
                () =>
                    {
                        foreach (var x in Enumerable.Range(0, 10))
                        {
                            queue.Add(x);
                        }
                    });

            i = 0;
            foreach (var element in queue.GetConsumingEnumerable(
                new CancellationTokenSource(3000).Token).Split(3))
            {
                Console.WriteLine($"Fetched chunk: {element.Count}");
                Console.WriteLine($"Fetched total: {i += element.Count}");
            }

            //Fetched chunk: 3
            //Fetched total: 3
            //Fetched chunk: 3
            //Fetched total: 6
            //Fetched chunk: 3
            //Fetched total: 9
        }

显然,如果元素少于块大小,则最后一个块将被“丢弃”。 有任何想法吗?

【问题讨论】:

  • 你想做什么?描述实际问题,而不是解决问题的尝试。 BlockingCollection不是用于流处理。为此有专门构建的库和类,例如 TPL 数据流或通道。 BatchBlock 将使用一行代码将传入的消息批量分成 N 个项目的批次。 ActionBlockTransformBlock 将使用 1 个或多个工作任务处理传入消息 LinkTo 将消息从一个块传递到另一个块而无需额外代码。几乎所有数据流块类型都有内置输入和输出缓冲区(如果适用)
  • 谢谢,我会看看那些。我仍然很好奇到底是什么导致了这个问题。 GetConsumingEnumerable 公开了一个 IEnumerable,我应该能够根据需要对其进行迭代。
  • 如果您使用 Chunk LINQ 运算符而不是 Split,问题是否仍然存在?

标签: c# producer-consumer blockingcollection


【解决方案1】:

我们应该打电话完成添加()通知方法GetConsumingEnumerable()没有更多的元素可以从它的生产者那里添加。

以下代码更改将解决您的问题并打印缺失的行。

Task.Run(() =>
              {
                 foreach (var x in Enumerable.Range(0, 10))
                 {
                    queue.Add(x);
                 }
                    queue.CompleteAdding(); // After executing the line, IsCompleted property of queue will be true.
              });

有关 BlockingCollection 的 GetConsumingEnumerable() 的更多信息,请参考this 链接。

【讨论】:

  • 只是为了添加有关行为的更多信息。取消将在 3 秒后发生,它将中止操作(意味着丢弃剩余的元素而不是返回它们(在我们的例子中是一个元素))。
猜你喜欢
  • 2023-01-30
  • 2019-09-24
  • 2014-05-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多