Skip to content

Add a performant groupWithin that works on chunks - #3775

Open
mmienko wants to merge 13 commits into
typelevel:mainfrom
mmienko:groupChunkWithin
Open

mmienko wants to merge 13 commits into
typelevel:mainfrom
mmienko:groupChunkWithin

Conversation

@mmienko

@mmienko mmienko commented Oct 2, 2026 •

Copy link
Copy Markdown

Add a performant groupWithin impl that works on chunks.

I experimented without other concurrency primitives like a Ref and Synchronous Queue, or a Ref of 2 Deferred, but they all obfuscated the main logic of this method and made it difficult to verify it's correct. SignallingRef had the perfect api, but was not performant in benchmarks. So, I created a trimmed down version called ConditionedRef which only pushes updates if they pass a predicate.

Comment thread core/shared/src/main/scala/fs2/Stream.scala
Comment thread core/shared/src/main/scala/fs2/concurrent/ConditionedRef.scala Outdated
Comment thread core/shared/src/main/scala/fs2/concurrent/ConditionedRef.scala
def modify[B](f: A => (A, B)): F[B] =
state.flatModify { s => // uncancellable to avoid losing wake-up signals
val (value, result) = f(s.value)
if (!s.waiters.exists(_.accepts(value)))

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.

Have you tried partitioning first to avoid two traversals of the list?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

We're dealing with a single consumer within groupWithin and the predicate is cheap vs. allocating new lists each time in partition where in practice the first if-clause should execute more often.

Arguably, this could be premature optimization (and I can run quick bechmarks to confirm) or impl should reflect spsc. Open to suggestions, but I'm starting to lean towards a scoped impl of this class for spsc rather than a more generic mpmc. Probably separate impls, ConditiedRef.mpmc, ConditiedRef.spsc (only this would be needed for PR), etc. will avoid confusion for future maintenance too.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

So there's a small hit when buffer size is small, but overall not very noticable. I'm fine with changing it to call partition first and have a single traversal.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

Ok, ran some more benchmarks and there's a drop in perf in the scenarios where there is no chunking (single elements), so I would prefer to keep it for my use case.

Comment thread core/shared/src/main/scala/fs2/Chunk.scala Outdated
Comment thread core/shared/src/main/scala/fs2/Chunk.scala Outdated
Comment thread core/shared/src/main/scala/fs2/Stream.scala
@mmienko

mmienko commented Oct 5, 2026

Copy link
Copy Markdown
Author

Benchmarks show that the new impl is improvement in all benchmarked situations, and especially when using a larger buffer size and/or chunks from upstream producer.

[info] # JMH version: 1.37
[info] # VM version: JDK 22.0.2, OpenJDK 64-Bit Server VM, 22.0.2+9-70
[info] # VM invoker: /Library/Java/JavaVirtualMachines/jdk-22.0.2.jdk/Contents/Home/bin/java
[info] # VM options: -Xms256m -Xmx256m

GroupWithinBenchmark
[info] Benchmark                                        (bufferSize)  (rangeLength)   Mode  Cnt     Score     Error  Units
[info] GroupWithinBenchmark.groupWithin                           16            100  thrpt   25  5045.189 ±  57.533  ops/s
[info] GroupWithinBenchmark.groupWithin                           16          10000  thrpt   25    66.093 ±   0.722  ops/s
[info] GroupWithinBenchmark.groupWithin                           16         100000  thrpt   25     6.685 ±   0.477  ops/s
[info] GroupWithinBenchmark.groupWithin                          256            100  thrpt   25  6436.653 ± 208.037  ops/s
[info] GroupWithinBenchmark.groupWithin                          256          10000  thrpt   25    92.835 ±   1.057  ops/s
[info] GroupWithinBenchmark.groupWithin                          256         100000  thrpt   25     7.856 ±   2.258  ops/s
[info] GroupWithinBenchmark.groupWithin                         4096            100  thrpt   25  6484.342 ± 125.325  ops/s
[info] GroupWithinBenchmark.groupWithin                         4096          10000  thrpt   25    94.488 ±   1.166  ops/s
[info] GroupWithinBenchmark.groupWithin                         4096         100000  thrpt   25     9.193 ±   0.162  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream            16            100  thrpt   25  4076.886 ±  75.363  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream            16          10000  thrpt   25    49.178 ±   1.052  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream            16         100000  thrpt   25     3.881 ±   1.209  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream           256            100  thrpt   25  5525.203 ± 222.971  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream           256          10000  thrpt   25    76.153 ±   1.327  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream           256         100000  thrpt   25     7.427 ±   0.143  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream          4096            100  thrpt   25  5627.852 ± 159.950  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream          4096          10000  thrpt   25    76.840 ±   1.502  ops/s
[info] GroupWithinBenchmark.groupWithinChunkedUpstream          4096         100000  thrpt   25     7.666 ±   0.144  ops/s

GroupChunksWithinBenchmark
[info] Benchmark                                                    (bufferSize)  (rangeLength)   Mode  Cnt      Score      Error  Units
[info] GroupChunksWithinBenchmark.groupChunksWithin                           16            100  thrpt   25   5922.350 ±   91.779  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithin                           16          10000  thrpt   25     71.620 ±    2.300  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithin                           16         100000  thrpt   25      7.434 ±    0.167  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithin                          256            100  thrpt   25   8108.692 ±  239.876  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithin                          256          10000  thrpt   25    110.113 ±   11.781  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithin                          256         100000  thrpt   25     10.584 ±    1.691  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithin                         4096            100  thrpt   25   7956.295 ±  175.919  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithin                         4096          10000  thrpt   25    109.900 ±    9.033  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithin                         4096         100000  thrpt   25      8.299 ±    2.447  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream            16            100  thrpt   25   7218.236 ±  124.838  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream            16          10000  thrpt   25     90.066 ±    8.578  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream            16         100000  thrpt   25      9.259 ±    0.140  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream           256            100  thrpt   25  11828.415 ± 2984.112  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream           256          10000  thrpt   25    181.671 ±   59.093  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream           256         100000  thrpt   25     25.281 ±    1.452  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream          4096            100  thrpt   25  14267.959 ± 1376.390  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream          4096          10000  thrpt   25    241.713 ±   83.001  ops/s
[info] GroupChunksWithinBenchmark.groupChunksWithinChunkedUpstream          4096         100000  thrpt   25     31.569 ±    0.761  ops/s

Based on this, we should swap the impl of groupWithin so that there is only one method.

Should we commit these to a text file under the benchmarks module?

@reardonj

reardonj commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

Should we commit these to a text file under the benchmarks module?

No. I don't see any other benchmarks with results committed.

@mmienko mmienko changed the title Add a performant groupChunksWithin Add a performant groupWithin that works on chunks Oct 5, 2026
Comment thread core/shared/src/main/scala/fs2/Chunk.scala

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants