Skip to content

optimize BlockingCollection<T>.GetConsumingEnumerable iterator - #69631

Closed
pedrobsaila wants to merge 4 commits into
dotnet:mainfrom
pedrobsaila:69320
Closed

pedrobsaila wants to merge 4 commits into
dotnet:mainfrom
pedrobsaila:69320

Conversation

@pedrobsaila

@pedrobsaila pedrobsaila commented May 20, 2022 •

Copy link
Copy Markdown
Contributor

Fixes #69320

I tried to search why this double-check had been done. It exists since the first RC of .NETCore 1.0. Could not find what version of .NET Framework introduced this change (just found that it exists in .NET 4.8)

@ghost ghost added area-System.Collections community-contribution Indicates that the PR has been added by a community member labels May 20, 2022
@ghost

ghost commented May 20, 2022

Copy link
Copy Markdown

Tagging subscribers to this area: @dotnet/area-system-collections
See info in area-owners.md if you want to be subscribed.

Issue Details

Fixes #69320

I tried to search why this double-check had been done. In history, it exists since .NETCore 1.0 RC. Could not find what version of .NET Framework introduced this change (just found that it exists in .NET 4.8)

Author: pedrobsaila
Assignees: -
Labels:

area-System.Collections

Milestone: -

@theodorzoulias

Copy link
Copy Markdown
Contributor

Hi @pedrobsaila! Please take a look at this recent comment, regarding the behavior of the BlockingCollection<T> in an edge case. Your PR as it stands now probably changes the current behavior of the collection.

cancellationToken.ThrowIfCancellationRequested();

//If the collection is completed then there is no need to wait.
if (IsCompleted)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should be replaced with Debug.Assert(IsCompleted);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should be replaced with Debug.Assert(IsCompleted);

IsCompleted can be changed from antoher thread , the assert would not be always true https://helixre8s23ayyeko0k025g8.blob.core.windows.net/dotnet-runtime-refs-pull-69631-merge-a83f55ae749f443aa3/System.Collections.Concurrent.Tests/1/console.b9f2d628.log?helixlogtype=result

@@ -675,11 +699,6 @@ private bool TryTakeWithNoTimeValidation([MaybeNullWhen(false)] out T item, int

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can the CheckDisposed(); call in line 673 also be removed in favor of an assert? It seems like IsCompleted is also making that check.

@eiriktsarpalis eiriktsarpalis left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm not sure if we could call this change an "optimization" unless there are demonstrable performance improvements in benchmarks.

@eiriktsarpalis eiriktsarpalis self-assigned this May 21, 2022
@pedrobsaila

pedrobsaila commented May 21, 2022 •

Copy link
Copy Markdown
Contributor Author

I'm not sure if we could call this change an "optimization" unless there are demonstrable performance improvements in benchmarks.

Here's the source code of the benchmark dotnet/performance@eb01ff1. I do benchmark for the worst case which is when collection is not complete. The opposite case will have the same results I guess, my change is not supposed to impact it

  • Test scenario AddRemoveFromDifferentThreads<T>.BlockingCollection :
    • Without double check :
AddRemoveFromDifferentThreads<Int32>.BlockingCollection: Job-RYUCKU(PowerPlanMode=00000000-0000-0000-0000-000000000000, Toolchain=CoreRun, InvocationCount=1, IterationTime=250.0000 ms, MaxIterationCount=20, MaxWarmupIterationCount=10, MinIterationCount=15, MinWarmupIterationCount=6, UnrollFactor=1, WarmupCount=-1) [Size=2000000]

BenchmarkDotNet=v0.13.1.1786-nightly, OS=Windows 10 (10.0.19044.1706/21H2/November2021Update)
Intel Core i7-10875H CPU 2.30GHz, 1 CPU, 16 logical and 8 physical cores
.NET SDK=7.0.100-preview.1.22110.4
  [Host]     : .NET 7.0.0 (7.0.22.7608), X64 RyuJIT
  Job-RYUCKU : .NET 7.0.0 (42.42.42.42424), X64 RyuJIT

PowerPlanMode=00000000-0000-0000-0000-000000000000  Toolchain=CoreRun  InvocationCount=1
IterationTime=250.0000 ms  MaxIterationCount=20  MaxWarmupIterationCount=10
MinIterationCount=15  MinWarmupIterationCount=6  UnrollFactor=1
WarmupCount=-1

|             Method |    Size |     Mean |    Error |   StdDev |   Median |      Min |      Max |      Gen 0 |     Gen 1 |     Gen 2 | Allocated |
|------------------- |-------- |---------:|---------:|---------:|---------:|---------:|---------:|-----------:|----------:|----------:|----------:|
| BlockingCollection | 2000000 | 400.9 ms | 15.39 ms | 17.72 ms | 406.5 ms | 367.3 ms | 422.7 ms | 30000.0000 | 1000.0000 | 1000.0000 | 260.14 MB |

AddRemoveFromDifferentThreads<String>.BlockingCollection: Job-RYUCKU(PowerPlanMode=00000000-0000-0000-0000-000000000000, Toolchain=CoreRun, InvocationCount=1, IterationTime=250.0000 ms, MaxIterationCount=20, MaxWarmupIterationCount=10, MinIterationCount=15, MinWarmupIterationCount=6, UnrollFactor=1, WarmupCount=-1) [Size=2000000]

BenchmarkDotNet=v0.13.1.1786-nightly, OS=Windows 10 (10.0.19044.1706/21H2/November2021Update)
Intel Core i7-10875H CPU 2.30GHz, 1 CPU, 16 logical and 8 physical cores
.NET SDK=7.0.100-preview.1.22110.4
  [Host]     : .NET 7.0.0 (7.0.22.7608), X64 RyuJIT
  Job-RYUCKU : .NET 7.0.0 (42.42.42.42424), X64 RyuJIT

PowerPlanMode=00000000-0000-0000-0000-000000000000  Toolchain=CoreRun  InvocationCount=1
IterationTime=250.0000 ms  MaxIterationCount=20  MaxWarmupIterationCount=10
MinIterationCount=15  MinWarmupIterationCount=6  UnrollFactor=1
WarmupCount=-1

|             Method |    Size |     Mean |    Error |   StdDev |   Median |      Min |      Max |      Gen 0 |     Gen 1 |     Gen 2 | Allocated |
|------------------- |-------- |---------:|---------:|---------:|---------:|---------:|---------:|-----------:|----------:|----------:|----------:|
| BlockingCollection | 2000000 | 452.8 ms | 20.90 ms | 23.23 ms | 453.4 ms | 411.5 ms | 509.4 ms | 31000.0000 | 2000.0000 | 2000.0000 | 276.14 MB |
  • With double check :
AddRemoveFromDifferentThreads<Int32>.BlockingCollection: Job-ZTJPHQ(PowerPlanMode=00000000-0000-0000-0000-000000000000, Toolchain=CoreRun, InvocationCount=1, IterationTime=250.0000 ms, MaxIterationCount=20, MaxWarmupIterationCount=10, MinIterationCount=15, MinWarmupIterationCount=6, UnrollFactor=1, WarmupCount=-1) [Size=2000000]

BenchmarkDotNet=v0.13.1.1786-nightly, OS=Windows 10 (10.0.19044.1706/21H2/November2021Update)
Intel Core i7-10875H CPU 2.30GHz, 1 CPU, 16 logical and 8 physical cores
.NET SDK=7.0.100-preview.1.22110.4
  [Host]     : .NET 7.0.0 (7.0.22.7608), X64 RyuJIT
  Job-ZTJPHQ : .NET 7.0.0 (42.42.42.42424), X64 RyuJIT

PowerPlanMode=00000000-0000-0000-0000-000000000000  Toolchain=CoreRun  InvocationCount=1
IterationTime=250.0000 ms  MaxIterationCount=20  MaxWarmupIterationCount=10
MinIterationCount=15  MinWarmupIterationCount=6  UnrollFactor=1
WarmupCount=-1

|             Method |    Size |     Mean |    Error |   StdDev |   Median |      Min |      Max |      Gen 0 |     Gen 1 |     Gen 2 | Allocated |
|------------------- |-------- |---------:|---------:|---------:|---------:|---------:|---------:|-----------:|----------:|----------:|----------:|
| BlockingCollection | 2000000 | 428.9 ms | 16.28 ms | 18.75 ms | 432.1 ms | 395.4 ms | 465.6 ms | 30000.0000 | 1000.0000 | 1000.0000 | 260.14 MB |

AddRemoveFromDifferentThreads<String>.BlockingCollection: Job-ZTJPHQ(PowerPlanMode=00000000-0000-0000-0000-000000000000, Toolchain=CoreRun, InvocationCount=1, IterationTime=250.0000 ms, MaxIterationCount=20, MaxWarmupIterationCount=10, MinIterationCount=15, MinWarmupIterationCount=6, UnrollFactor=1, WarmupCount=-1) [Size=2000000]

BenchmarkDotNet=v0.13.1.1786-nightly, OS=Windows 10 (10.0.19044.1706/21H2/November2021Update)
Intel Core i7-10875H CPU 2.30GHz, 1 CPU, 16 logical and 8 physical cores
.NET SDK=7.0.100-preview.1.22110.4
  [Host]     : .NET 7.0.0 (7.0.22.7608), X64 RyuJIT
  Job-ZTJPHQ : .NET 7.0.0 (42.42.42.42424), X64 RyuJIT

PowerPlanMode=00000000-0000-0000-0000-000000000000  Toolchain=CoreRun  InvocationCount=1
IterationTime=250.0000 ms  MaxIterationCount=20  MaxWarmupIterationCount=10
MinIterationCount=15  MinWarmupIterationCount=6  UnrollFactor=1
WarmupCount=-1

|             Method |    Size |     Mean |    Error |   StdDev |   Median |      Min |      Max |      Gen 0 |     Gen 1 |     Gen 2 | Allocated |
|------------------- |-------- |---------:|---------:|---------:|---------:|---------:|---------:|-----------:|----------:|----------:|----------:|
| BlockingCollection | 2000000 | 448.5 ms | 17.31 ms | 19.24 ms | 444.8 ms | 413.2 ms | 486.6 ms | 31000.0000 | 2000.0000 | 2000.0000 | 276.14 MB |
  • Test scenario AddRemoveFromSameThreads<T>.BlockingCollection :
    • Without double check :
AddRemoveFromSameThreads<Int32>.BlockingCollection: Job-OTUVIK(PowerPlanMode=00000000-0000-0000-0000-000000000000, Toolchain=CoreRun, InvocationCount=1, IterationTime=250.0000 ms, MaxIterationCount=20, MaxWarmupIterationCount=10, MinIterationCount=15, MinWarmupIterationCount=6, UnrollFactor=1, WarmupCount=-1) [Size=2000000]

BenchmarkDotNet=v0.13.1.1786-nightly, OS=Windows 10 (10.0.19044.1706/21H2/November2021Update)
Intel Core i7-10875H CPU 2.30GHz, 1 CPU, 16 logical and 8 physical cores
.NET SDK=7.0.100-preview.1.22110.4
  [Host]     : .NET 7.0.0 (7.0.22.7608), X64 RyuJIT
  Job-OTUVIK : .NET 7.0.0 (42.42.42.42424), X64 RyuJIT

PowerPlanMode=00000000-0000-0000-0000-000000000000  Toolchain=CoreRun  InvocationCount=1
IterationTime=250.0000 ms  MaxIterationCount=20  MaxWarmupIterationCount=10
MinIterationCount=15  MinWarmupIterationCount=6  UnrollFactor=1
WarmupCount=-1

|             Method |    Size |     Mean |    Error |   StdDev |   Median |      Min |      Max |      Gen 0 |     Gen 1 |     Gen 2 | Allocated |
|------------------- |-------- |---------:|---------:|---------:|---------:|---------:|---------:|-----------:|----------:|----------:|----------:|
| BlockingCollection | 2000000 | 498.4 ms | 37.05 ms | 42.66 ms | 496.4 ms | 440.9 ms | 580.3 ms | 31000.0000 | 2000.0000 | 2000.0000 | 276.15 MB |

AddRemoveFromSameThreads<String>.BlockingCollection: Job-OTUVIK(PowerPlanMode=00000000-0000-0000-0000-000000000000, Toolchain=CoreRun, InvocationCount=1, IterationTime=250.0000 ms, MaxIterationCount=20, MaxWarmupIterationCount=10, MinIterationCount=15, MinWarmupIterationCount=6, UnrollFactor=1, WarmupCount=-1) [Size=2000000]

BenchmarkDotNet=v0.13.1.1786-nightly, OS=Windows 10 (10.0.19044.1706/21H2/November2021Update)
Intel Core i7-10875H CPU 2.30GHz, 1 CPU, 16 logical and 8 physical cores
.NET SDK=7.0.100-preview.1.22110.4
  [Host]     : .NET 7.0.0 (7.0.22.7608), X64 RyuJIT
  Job-OTUVIK : .NET 7.0.0 (42.42.42.42424), X64 RyuJIT

PowerPlanMode=00000000-0000-0000-0000-000000000000  Toolchain=CoreRun  InvocationCount=1
IterationTime=250.0000 ms  MaxIterationCount=20  MaxWarmupIterationCount=10
MinIterationCount=15  MinWarmupIterationCount=6  UnrollFactor=1
WarmupCount=-1

|             Method |    Size |     Mean |    Error |   StdDev |   Median |      Min |      Max |      Gen 0 |     Gen 1 |     Gen 2 | Allocated |
|------------------- |-------- |---------:|---------:|---------:|---------:|---------:|---------:|-----------:|----------:|----------:|----------:|
| BlockingCollection | 2000000 | 546.2 ms | 28.08 ms | 32.34 ms | 541.6 ms | 501.0 ms | 620.9 ms | 31000.0000 | 2000.0000 | 2000.0000 | 308.15 MB |
  • With double check :
AddRemoveFromSameThreads<Int32>.BlockingCollection: Job-BWWGOB(PowerPlanMode=00000000-0000-0000-0000-000000000000, Toolchain=CoreRun, InvocationCount=1, IterationTime=250.0000 ms, MaxIterationCount=20, MaxWarmupIterationCount=10, MinIterationCount=15, MinWarmupIterationCount=6, UnrollFactor=1, WarmupCount=-1) [Size=2000000]

BenchmarkDotNet=v0.13.1.1786-nightly, OS=Windows 10 (10.0.19044.1706/21H2/November2021Update)
Intel Core i7-10875H CPU 2.30GHz, 1 CPU, 16 logical and 8 physical cores
.NET SDK=7.0.100-preview.1.22110.4
  [Host]     : .NET 7.0.0 (7.0.22.7608), X64 RyuJIT
  Job-BWWGOB : .NET 7.0.0 (42.42.42.42424), X64 RyuJIT

PowerPlanMode=00000000-0000-0000-0000-000000000000  Toolchain=CoreRun  InvocationCount=1
IterationTime=250.0000 ms  MaxIterationCount=20  MaxWarmupIterationCount=10
MinIterationCount=15  MinWarmupIterationCount=6  UnrollFactor=1
WarmupCount=-1

|             Method |    Size |     Mean |    Error |   StdDev |   Median |      Min |      Max |      Gen 0 |     Gen 1 |     Gen 2 | Allocated |
|------------------- |-------- |---------:|---------:|---------:|---------:|---------:|---------:|-----------:|----------:|----------:|----------:|
| BlockingCollection | 2000000 | 516.9 ms | 36.66 ms | 42.22 ms | 523.8 ms | 445.6 ms | 610.5 ms | 31000.0000 | 2000.0000 | 2000.0000 | 276.15 MB |

AddRemoveFromSameThreads<String>.BlockingCollection: Job-BWWGOB(PowerPlanMode=00000000-0000-0000-0000-000000000000, Toolchain=CoreRun, InvocationCount=1, IterationTime=250.0000 ms, MaxIterationCount=20, MaxWarmupIterationCount=10, MinIterationCount=15, MinWarmupIterationCount=6, UnrollFactor=1, WarmupCount=-1) [Size=2000000]

BenchmarkDotNet=v0.13.1.1786-nightly, OS=Windows 10 (10.0.19044.1706/21H2/November2021Update)
Intel Core i7-10875H CPU 2.30GHz, 1 CPU, 16 logical and 8 physical cores
.NET SDK=7.0.100-preview.1.22110.4
  [Host]     : .NET 7.0.0 (7.0.22.7608), X64 RyuJIT
  Job-BWWGOB : .NET 7.0.0 (42.42.42.42424), X64 RyuJIT

PowerPlanMode=00000000-0000-0000-0000-000000000000  Toolchain=CoreRun  InvocationCount=1
IterationTime=250.0000 ms  MaxIterationCount=20  MaxWarmupIterationCount=10
MinIterationCount=15  MinWarmupIterationCount=6  UnrollFactor=1
WarmupCount=-1

|             Method |    Size |     Mean |    Error |   StdDev |   Median |      Min |      Max |      Gen 0 |     Gen 1 |     Gen 2 | Allocated |
|------------------- |-------- |---------:|---------:|---------:|---------:|---------:|---------:|-----------:|----------:|----------:|----------:|
| BlockingCollection | 2000000 | 566.6 ms | 41.47 ms | 47.76 ms | 560.5 ms | 511.2 ms | 672.3 ms | 31000.0000 | 2000.0000 | 2000.0000 | 308.15 MB |

item = default(T);
return false;
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@pedrobsaila this alters the behavior of the Take method in case the collection is completed and the cancellationToken is canceled. The existing behavior is to throw an OperationCanceledException. The new behavior will be to return false. Altering the behavior of the method is unlikely to be desirable. A demonstration of the current behavior can be found in this comment.

{
ValidateTimeout(timeout);

//If the collection is completed then there is no need remove an item.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Maybe:

Suggested change
//If the collection is completed then there is no need remove an item.
// If the collection is completed then there is no need remove an item.

@eiriktsarpalis

Copy link
Copy Markdown
Member

Without double check

Can you clarify what is meant by a "double check" in this case? Judging by the diff it seems that you have moved the check out of TryTakeWithNoTimeValidation and into the methods calling it.

@pedrobsaila

Copy link
Copy Markdown
Contributor Author

Without double check

Is the code in pedrobsaila:69320 branch

With double check

Is the code in main branch

@eiriktsarpalis

Copy link
Copy Markdown
Member

@pedrobsaila I get that, but I fail to see which code path avoids checking IsCompleted twice. At short glance the diff seems to simply be moving the check out of the TryTakeWithNoTimeValidation method and into the methods that call it.

@pedrobsaila

pedrobsaila commented Jun 9, 2022 •

Copy link
Copy Markdown
Contributor Author

@pedrobsaila I get that, but I fail to see which code path avoids checking IsCompleted twice. At short glance the diff seems to simply be moving the check out of the TryTakeWithNoTimeValidation method and into the methods that call it.

It is the IEnumerable<T> GetConsumingEnumerable(CancellationToken cancellationToken) and IEnumerable<T> GetConsumingEnumerable() methods :

Call stack for : IEnumerable<T> GetConsumingEnumerable(CancellationToken cancellationToken) /*calls IsCompleted in both branches)*/
At bool TryTakeWithNoTimeValidation([MaybeNullWhen(false)] out T item, int millisecondsTimeout,
CancellationToken cancellationToken, CancellationTokenSource? combinedTokenSource) /*calls IsCompleted in main, but does not in pedrobsaila:69320*/

Call stack for : IEnumerable<T> GetConsumingEnumerable()
At IEnumerable<T> GetConsumingEnumerable(CancellationToken cancellationToken) /*calls IsCompleted in both branches)*/
At bool TryTakeWithNoTimeValidation([MaybeNullWhen(false)] out T item, int millisecondsTimeout,
CancellationToken cancellationToken, CancellationTokenSource? combinedTokenSource) /*calls IsCompleted in main, but does not in pedrobsaila:69320*/

@eiriktsarpalis

Copy link
Copy Markdown
Member

I see, so it basically only concerns this callsite?

while (!IsCompleted)
{
T? item;
if (TryTakeWithNoTimeValidation(out item, Timeout.Infinite, cancellationToken, linkedTokenSource))
{
yield return item;
}
}

Not necessarily convinced shaving a few percentage points is worth the potential risk of accidental regression. If anything, BlockingCollection has been superseded by more modern patterns such as System.Threading.Channels so we should encourage users to adopt that instead.

@pedrobsaila

Copy link
Copy Markdown
Contributor Author

I see, so it basically only concerns this callsite?

while (!IsCompleted)
{
T? item;
if (TryTakeWithNoTimeValidation(out item, Timeout.Infinite, cancellationToken, linkedTokenSource))
{
yield return item;
}
}

Exactly

Not necessarily convinced shaving a few percentage points is worth the potential risk of accidental regression. If anything, BlockingCollection has been superseded by more modern patterns such as System.Threading.Channels so we should encourage users to adopt that instead.

I understand. There is not a substantial improvement and if we take into consideration the error, maybe there is any of all. I just wanted to try doing something related to perf for the first time. You can close it.

@theodorzoulias

Copy link
Copy Markdown
Contributor

I understand. There is not a substantial improvement and if we take into consideration the error, maybe there is any of all. I just wanted to try doing something related to perf for the first time. You can close it.

@pedrobsaila when I made the initial proposal, my idea was about a localized change inside the GetConsumingEnumerable method, that would improve slightly the performance of this method. I was not thinking about moving code around in half a dozen of places inside the BlockingCollection<T> class! Currently I am not sure that this PR preserves the existing cancellation behavior of the class. Apparently there are no tests in place that will fail if this behavior has changed. So, for what it's worth, I would not advocate about changing this class, unless the changes are few and localized, and tests for the behavior of the affected methods are included in the PR. 🙂

@theodorzoulias

theodorzoulias commented Jun 10, 2022 •

Copy link
Copy Markdown
Contributor

If anything, BlockingCollection has been superseded by more modern patterns such as System.Threading.Channels so we should encourage users to adopt that instead.

@eiriktsarpalis my understanding is that the Channels are intended for asynchronous scenarios. Using them is synchronous scenarios will require blocking on asynchronous methods that return ValueTasks, that require converting these to Tasks with AsTask, that will result in an allocation on each blocking operation. Which is something undesirable in general. So, unless my understanding is wrong, the BlockingCollection<T> has still a place in modern concurrent programming.

@pedrobsaila

pedrobsaila commented Jun 10, 2022 •

Copy link
Copy Markdown
Contributor Author

@pedrobsaila when I made the initial proposal

I think you are talking about this :

while (true)
{
    T? item;
    if (TryTakeWithNoTimeValidation(out item, Timeout.Infinite, cancellationToken, linkedTokenSource))
    {
        yield return item;
    }
    else
    {
        break;
    }
}

I did not follow the initial proposal, because (correct me if I'm wrong) your proposal assumes that when there is no element to take, GetConsumingEnumerable should end, whereas the actual code retry taking another element untill adding is complete

@theodorzoulias

Copy link
Copy Markdown
Contributor

@pedrobsaila yes, this is the change I had in mind. Surprisingly the TryTakeWithNoTimeValidation is a blocking method. It will block indefinitely (Timeout.Infinite) until either an item is available, or the collection is completed. The name of this method is confusing, because in general the Try-prefixed methods in the standard libraries return immediately and don't block.

@eiriktsarpalis

Copy link
Copy Markdown
Member

@pedrobsaila in that case, I'm going to close this PR. Thank you for the contribution!

@stephentoub

Copy link
Copy Markdown
Member

@eiriktsarpalis, the original issue should be closed as well then?

@eiriktsarpalis

Copy link
Copy Markdown
Member

Per @theodorzoulias's comment in #69631 (comment) the original issue was proposing a different change, although I haven't validated the feasibility or usefulness of the proposal.

@pedrobsaila

pedrobsaila commented Jun 13, 2022 •

Copy link
Copy Markdown
Contributor Author

Per @theodorzoulias's comment in #69631 (comment) the original issue was proposing a different change, although I haven't validated the feasibility or usefulness of the proposal.

The initial proposal changes the behaviour of GetConsumingEnumerable because it makes it non-blocking when there's no element to take from collection. The main code blocks until IsCompleted.

@theodorzoulias

Copy link
Copy Markdown
Contributor

The initial proposal changes the behaviour of GetConsumingEnumerable because it makes it non-blocking when there's no element to take from collection. The main code blocks until IsCompleted.

@pedrobsaila have you confirmed this experimentally, or you infer it by reading the code?

@pedrobsaila

pedrobsaila commented Jun 14, 2022 •

Copy link
Copy Markdown
Contributor Author

@pedrobsaila have you confirmed this experimentally, or you infer it by reading the code?

yes I experimented it, try :

BlockingCollection<int> blockingCollection = new BlockingCollection<int>();
blockingCollection.Add(2);
foreach (int i in blockingCollection.GetConsumingEnumerable())
{
      Debug.WriteLine(i);
}

@theodorzoulias

theodorzoulias commented Jun 14, 2022 •

Copy link
Copy Markdown
Contributor

@pedrobsaila please take a look this online demo. I implemented the GetConsumingEnumerable2 extension method as shown in my original proposal, using reflection to read the private field _consumersCancellationTokenSource and invoke the private method TryTakeWithNoTimeValidation. Then I ran the code sample of your previous comment, replacing the GetConsumingEnumerable with GetConsumingEnumerable2. The resulting behavior is that the current thread is blocked indefinitely after writing 2 in the console. It's exactly the same behavior with the GetConsumingEnumerable method. Could you clarify what difference you observed in your experiments?

@pedrobsaila

pedrobsaila commented Jun 14, 2022 •

Copy link
Copy Markdown
Contributor Author

@pedrobsaila please take a look this online demo. I implemented the GetConsumingEnumerable2 extension method as shown in my original proposal, using reflection to read the private field _consumersCancellationTokenSource and invoke the private method TryTakeWithNoTimeValidation. Then I ran the code sample of your previous comment, replacing the GetConsumingEnumerable with GetConsumingEnumerable2. The resulting behavior is that the current thread is blocked indefinitely after writing 2 in the console. It's exactly the same behavior with the GetConsumingEnumerable method. Could you clarify what difference you observed in your experiments?

Yes the test blocks in your proposal also (sorry for that false alarm). It's because we block also when collection is empty with _occupiedNodes semaphore. I think your proposal is also good to go. I will try to make another PR with your own proposal with performance results

@theodorzoulias

Copy link
Copy Markdown
Contributor

@pedrobsaila be aware than my proposal changes the behavior of the GetConsumingEnumerable method, when an already canceled CancellationToken is passed as argument. My theory is that it can be fixed by adding an if (IsCompleted) yield break; at the very top of the method, but I haven't verified it experimentally. I am also guessing that there are no tests in place that verify the behavior of the GetConsumingEnumerable method in this edge case, which is a bit unnerving TBH!

@ghost ghost locked as resolved and limited conversation to collaborators Jul 15, 2022
@pedrobsaila
pedrobsaila deleted the 69320 branch November 26, 2022 16:06
Sign up for free to subscribe to this conversation on GitHub. Already have an account? Sign in.

Labels

area-System.Collections community-contribution Indicates that the PR has been added by a community member

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Can the BlockingCollection<T>.GetConsumingEnumerable iterator be optimized?

5 participants