This is an automated email from the ASF dual-hosted git repository.
He-Pin pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko.git
The following commit(s) were added to refs/heads/main by this push:
new fe166c73ee feat: add maxTotalBufferSize to MergeHub for aggregate
queue bound (#3392)
fe166c73ee is described below
commit fe166c73ee88b1864f0bfa9df4eeadbb1a69ef52
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Tue Jul 28 18:22:01 2026 +0800
feat: add maxTotalBufferSize to MergeHub for aggregate queue bound (#3392)
* feat: add maxTotalBufferSize to MergeHub for aggregate queue bound
Motivation:
MergeHub's central queue is unbounded. In dynamic producer scenarios
where producers continuously materialize while downstream is stalled,
the queue grows without limit leading to OOM (issue #3248).
Modification:
Add maxTotalBufferSize parameter to MergeHub (scaladsl + javadsl).
Track buffered element count with AtomicInteger. New producers are
cancelled at registration when the count meets or exceeds the threshold.
Default 0 preserves existing unlimited behavior.
Result:
Users can bound aggregate MergeHub queue growth via admission control,
preventing OOM in dynamic producer churn scenarios.
Tests:
- sbt "stream-tests / Test / testOnly
org.apache.pekko.stream.scaladsl.HubSpec -- -z MergeHub"
17 tests passed (including 2 new directional tests)
- sbt "stream / mimaReportBinaryIssues" passed
References:
Fixes #3248
* fix: address CR comments for maxTotalBufferSize PR
Motivation:
Address review comments from pjfanning on PR #3392 and improve
test determinism and documentation clarity.
Modification:
- Add @since 2.0.0 to all 4 new public methods (scaladsl + javadsl)
- Fix typo "produces" -> "producers" in javadsl withDraining
- Remove misleading "Default value is 16" from two-arg overloads
- Rewrite overshoot docs to clarify soft-bound admission semantics
- Add code comment explaining best-effort admission check
- Redesign tests with expectNext synchronization barrier to
eliminate race condition in admission rejection assertions
- Add new test: already admitted producers continue when full
Result:
CR comments resolved, tests are deterministic, docs are accurate.
Tests:
- sbt "stream-tests / Test / testOnly
org.apache.pekko.stream.scaladsl.HubSpec -- -z MergeHub"
18 tests passed
- scalafmt --mode diff-ref=origin/main — clean
- git diff --check — clean
References:
Refs #3248
---
.../org/apache/pekko/stream/scaladsl/HubSpec.scala | 97 ++++++++++++++++++++++
.../org/apache/pekko/stream/javadsl/Hub.scala | 64 ++++++++++++++
.../org/apache/pekko/stream/scaladsl/Hub.scala | 66 ++++++++++++++-
3 files changed, 223 insertions(+), 4 deletions(-)
diff --git
a/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/HubSpec.scala
b/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/HubSpec.scala
index deaebb4d51..f8d445d106 100644
--- a/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/HubSpec.scala
+++ b/stream-tests/src/test/scala/org/apache/pekko/stream/scaladsl/HubSpec.scala
@@ -153,6 +153,103 @@ class HubSpec extends StreamSpec {
sub.cancel()
}
+ "cancel new producers when maxTotalBufferSize is exceeded" in {
+ val downstream = TestSubscriber.manualProbe[Int]()
+ // perProducerBufferSize=4, maxTotalBufferSize=3
+ val sink =
Sink.fromSubscriber(downstream).runWith(MergeHub.source[Int](4, 3))
+
+ val sub = downstream.expectSubscription()
+
+ // Producer connects and pushes 4 elements
+ val upstream1 = TestPublisher.probe[Int]()
+ Source.fromPublisher(upstream1).runWith(sink)
+ for (i <- 1 to 4) upstream1.sendNext(i)
+
+ // Consume 1 element: proves all 4 were enqueued (counter was 4, now 3 =
threshold)
+ sub.request(1)
+ downstream.expectNext(1)
+
+ // Buffer count (3) >= maxTotalBufferSize (3): new producer is rejected
+ val upstream2 = TestPublisher.probe[Int]()
+ Source.fromPublisher(upstream2).runWith(sink)
+ upstream2.expectCancellation()
+
+ sub.cancel()
+ }
+
+ "allow new producers after elements are consumed from the aggregate
buffer" in {
+ val downstream = TestSubscriber.manualProbe[Int]()
+ // perProducerBufferSize=4, maxTotalBufferSize=3
+ val sink =
Sink.fromSubscriber(downstream).runWith(MergeHub.source[Int](4, 3))
+
+ val sub = downstream.expectSubscription()
+
+ // First producer connects and pushes 4 elements
+ val upstream1 = TestPublisher.probe[Int]()
+ Source.fromPublisher(upstream1).runWith(sink)
+ for (i <- 1 to 4) upstream1.sendNext(i)
+
+ // Consume 1 to prove all 4 were enqueued (counter was 4, now 3 =
threshold)
+ sub.request(1)
+ downstream.expectNext(1)
+
+ // Second producer is cancelled because buffer count (3) >= threshold (3)
+ val upstream2 = TestPublisher.probe[Int]()
+ Source.fromPublisher(upstream2).runWith(sink)
+ upstream2.expectCancellation()
+
+ // Consume all remaining elements, freeing the buffer (counter drops to
0 < 3)
+ sub.request(3)
+ downstream.expectNext(2)
+ downstream.expectNext(3)
+ downstream.expectNext(4)
+
+ // Now a new producer should be accepted and its element delivered
+ val upstream3 = TestPublisher.probe[Int]()
+ Source.fromPublisher(upstream3).runWith(sink)
+ upstream3.sendNext(5)
+ sub.request(1)
+ downstream.expectNext(5)
+
+ sub.cancel()
+ }
+
+ "allow already admitted producers to continue when buffer is full" in {
+ val downstream = TestSubscriber.manualProbe[Int]()
+ // perProducerBufferSize=4, maxTotalBufferSize=3
+ val sink =
Sink.fromSubscriber(downstream).runWith(MergeHub.source[Int](4, 3))
+
+ val sub = downstream.expectSubscription()
+
+ // First producer connects and pushes 4 elements (counter = 4 >
threshold 3)
+ val upstream1 = TestPublisher.probe[Int]()
+ Source.fromPublisher(upstream1).runWith(sink)
+ for (i <- 1 to 4) upstream1.sendNext(i)
+
+ // Consume 1 to prove elements were enqueued (counter = 3 = threshold)
+ sub.request(1)
+ downstream.expectNext(1)
+
+ // New producer is rejected
+ val upstream2 = TestPublisher.probe[Int]()
+ Source.fromPublisher(upstream2).runWith(sink)
+ upstream2.expectCancellation()
+
+ // But the already admitted producer can still push (it has remaining
demand)
+ // Consume all remaining to make room and generate demand signals
+ sub.request(3)
+ downstream.expectNext(2)
+ downstream.expectNext(3)
+ downstream.expectNext(4)
+
+ // First producer received demand signal and can push more
+ upstream1.sendNext(5)
+ sub.request(1)
+ downstream.expectNext(5)
+
+ sub.cancel()
+ }
+
"work with long streams" in {
val (sink, result) =
MergeHub.source[Int](16).take(2000).toMat(Sink.seq)(Keep.both).run()
Source(1 to 1000).runWith(sink)
diff --git a/stream/src/main/scala/org/apache/pekko/stream/javadsl/Hub.scala
b/stream/src/main/scala/org/apache/pekko/stream/javadsl/Hub.scala
index c6595b1f27..65aa2ce881 100644
--- a/stream/src/main/scala/org/apache/pekko/stream/javadsl/Hub.scala
+++ b/stream/src/main/scala/org/apache/pekko/stream/javadsl/Hub.scala
@@ -69,6 +69,35 @@ object MergeHub {
pekko.stream.scaladsl.MergeHub.source[T](perProducerBufferSize).mapMaterializedValue(_.asJava[T]).asJava
}
+ /**
+ * Creates a [[Source]] that emits elements merged from a dynamic set of
producers. After the [[Source]] returned
+ * by this method is materialized, it returns a [[Sink]] as a materialized
value. This [[Sink]] can be materialized
+ * arbitrary many times and each of the materializations will feed the
elements into the original [[Source]].
+ *
+ * Every new materialization of the [[Source]] results in a new, independent
hub, which materializes to its own
+ * [[Sink]] for feeding that materialization.
+ *
+ * Completed or failed [[Sink]]s are simply removed. Once the [[Source]] is
cancelled, the Hub is considered closed
+ * and any new producers using the [[Sink]] will be cancelled.
+ *
+ * @param clazz Type of elements this hub emits and consumes
+ * @param perProducerBufferSize Buffer space used per producer.
+ * @param maxTotalBufferSize Admission threshold for the total number of
elements buffered across all producers.
+ * New producers are cancelled at registration when the buffered element
count meets or exceeds this value.
+ * Already admitted producers are unaffected and may continue to push up
to their per-producer buffer.
+ * Use 0 for unlimited (default behavior).
+ * @since 2.0.0
+ */
+ def of[T](
+ @nowarn("msg=never used") clazz: Class[T],
+ perProducerBufferSize: Int,
+ maxTotalBufferSize: Int): Source[T, Sink[T, NotUsed]] = {
+ pekko.stream.scaladsl.MergeHub
+ .source[T](perProducerBufferSize, maxTotalBufferSize)
+ .mapMaterializedValue(_.asJava[T])
+ .asJava
+ }
+
/**
* Creates a [[Source]] that emits elements merged from a dynamic set of
producers. After the [[Source]] returned
* by this method is materialized, it returns a [[Sink]] as a materialized
value. This [[Sink]] can be materialized
@@ -98,6 +127,41 @@ object MergeHub {
.asJava
}
+ /**
+ * Creates a [[Source]] that emits elements merged from a dynamic set of
producers. After the [[Source]] returned
+ * by this method is materialized, it returns a [[Sink]] as a materialized
value. This [[Sink]] can be materialized
+ * arbitrarily many times and each of the materializations will feed the
elements into the original [[Source]].
+ *
+ * Every new materialization of the [[Source]] results in a new, independent
hub, which materializes to its own
+ * [[Sink]] for feeding that materialization.
+ *
+ * Completed or failed [[Sink]]s are simply removed. Once the [[Source]] is
cancelled, the Hub is considered closed
+ * and any new producers using the [[Sink]] will be cancelled.
+ *
+ * The materialized [[DrainingControl]] can be used to drain the Hub: any
new producers using the [[Sink]] will be cancelled
+ * and the Hub will be closed completing the [[Source]] as soon as all
currently connected producers complete.
+ *
+ * @param clazz Type of elements this hub emits and consumes
+ * @param perProducerBufferSize Buffer space used per producer.
+ * @param maxTotalBufferSize Admission threshold for the total number of
elements buffered across all producers.
+ * New producers are cancelled at registration when the buffered element
count meets or exceeds this value.
+ * Already admitted producers are unaffected and may continue to push up
to their per-producer buffer.
+ * Use 0 for unlimited (default behavior).
+ * @since 2.0.0
+ */
+ def withDraining[T](
+ @nowarn("msg=never used") clazz: Class[T],
+ perProducerBufferSize: Int,
+ maxTotalBufferSize: Int): Source[T, pekko.japi.Pair[Sink[T, NotUsed],
DrainingControl]] = {
+ pekko.stream.scaladsl.MergeHub
+ .sourceWithDraining[T](perProducerBufferSize, maxTotalBufferSize)
+ .mapMaterializedValue {
+ case (sink, draining) =>
+ pekko.japi.Pair(sink.asJava[T], new DrainingControlImpl(draining):
DrainingControl)
+ }
+ .asJava
+ }
+
/**
* Creates a [[Source]] that emits elements merged from a dynamic set of
producers. After the [[Source]] returned
* by this method is materialized, it returns a [[Sink]] as a materialized
value. This [[Sink]] can be materialized
diff --git a/stream/src/main/scala/org/apache/pekko/stream/scaladsl/Hub.scala
b/stream/src/main/scala/org/apache/pekko/stream/scaladsl/Hub.scala
index eee1129ad7..c95919a5ea 100644
--- a/stream/src/main/scala/org/apache/pekko/stream/scaladsl/Hub.scala
+++ b/stream/src/main/scala/org/apache/pekko/stream/scaladsl/Hub.scala
@@ -15,8 +15,7 @@ package org.apache.pekko.stream.scaladsl
import java.util
import java.util.concurrent.ConcurrentHashMap
-import java.util.concurrent.atomic.{ AtomicLong, AtomicReference }
-import java.util.concurrent.atomic.AtomicInteger
+import java.util.concurrent.atomic.{ AtomicInteger, AtomicLong,
AtomicReference }
import java.util.concurrent.atomic.AtomicReferenceArray
import scala.annotation.tailrec
@@ -78,6 +77,27 @@ object MergeHub {
def source[T](perProducerBufferSize: Int): Source[T, Sink[T, NotUsed]] =
Source.fromGraph(new MergeHub[T](perProducerBufferSize,
false)).mapMaterializedValue(_._1)
+ /**
+ * Creates a [[Source]] that emits elements merged from a dynamic set of
producers. After the [[Source]] returned
+ * by this method is materialized, it returns a [[Sink]] as a materialized
value. This [[Sink]] can be materialized
+ * arbitrary many times and each of the materializations will feed the
elements into the original [[Source]].
+ *
+ * Every new materialization of the [[Source]] results in a new, independent
hub, which materializes to its own
+ * [[Sink]] for feeding that materialization.
+ *
+ * Completed or failed [[Sink]]s are simply removed. Once the [[Source]] is
cancelled, the Hub is considered closed
+ * and any new producers using the [[Sink]] will be cancelled.
+ *
+ * @param perProducerBufferSize Buffer space used per producer.
+ * @param maxTotalBufferSize Admission threshold for the total number of
elements buffered across all producers.
+ * New producers are cancelled at registration when the buffered element
count meets or exceeds this value.
+ * Already admitted producers are unaffected and may continue to push up
to their per-producer buffer.
+ * Use 0 for unlimited (default behavior).
+ * @since 2.0.0
+ */
+ def source[T](perProducerBufferSize: Int, maxTotalBufferSize: Int):
Source[T, Sink[T, NotUsed]] =
+ Source.fromGraph(new MergeHub[T](perProducerBufferSize, false,
maxTotalBufferSize)).mapMaterializedValue(_._1)
+
/**
* Creates a [[Source]] that emits elements merged from a dynamic set of
producers. After the [[Source]] returned
* by this method is materialized, it returns a [[Sink]] as a materialized
value. This [[Sink]] can be materialized
@@ -97,6 +117,32 @@ object MergeHub {
def sourceWithDraining[T](perProducerBufferSize: Int): Source[T, (Sink[T,
NotUsed], DrainingControl)] =
Source.fromGraph(new MergeHub[T](perProducerBufferSize, true))
+ /**
+ * Creates a [[Source]] that emits elements merged from a dynamic set of
producers. After the [[Source]] returned
+ * by this method is materialized, it returns a [[Sink]] as a materialized
value. This [[Sink]] can be materialized
+ * arbitrary many times and each of the materializations will feed the
elements into the original [[Source]].
+ *
+ * Every new materialization of the [[Source]] results in a new, independent
hub, which materializes to its own
+ * [[Sink]] for feeding that materialization.
+ *
+ * Completed or failed [[Sink]]s are simply removed. Once the [[Source]] is
cancelled, the Hub is considered closed
+ * and any new producers using the [[Sink]] will be cancelled.
+ *
+ * The materialized [[DrainingControl]] can be used to drain the Hub: any
new producers using the [[Sink]] will be cancelled
+ * and the Hub will be closed completing the [[Source]] as soon as all
currently connected producers complete.
+ *
+ * @param perProducerBufferSize Buffer space used per producer.
+ * @param maxTotalBufferSize Admission threshold for the total number of
elements buffered across all producers.
+ * New producers are cancelled at registration when the buffered element
count meets or exceeds this value.
+ * Already admitted producers are unaffected and may continue to push up
to their per-producer buffer.
+ * Use 0 for unlimited (default behavior).
+ * @since 2.0.0
+ */
+ def sourceWithDraining[T](
+ perProducerBufferSize: Int,
+ maxTotalBufferSize: Int): Source[T, (Sink[T, NotUsed], DrainingControl)]
=
+ Source.fromGraph(new MergeHub[T](perProducerBufferSize, true,
maxTotalBufferSize))
+
/**
* Creates a [[Source]] that emits elements merged from a dynamic set of
producers. After the [[Source]] returned
* by this method is materialized, it returns a [[Sink]] as a materialized
value. This [[Sink]] can be materialized
@@ -140,9 +186,13 @@ private[pekko] final class
MergeHubDrainingControlImpl(drainAction: () => Unit)
}
}
-private[pekko] class MergeHub[T](perProducerBufferSize: Int, drainingEnabled:
Boolean = false)
+private[pekko] class MergeHub[T](
+ perProducerBufferSize: Int,
+ drainingEnabled: Boolean = false,
+ maxTotalBufferSize: Int = 0)
extends GraphStageWithMaterializedValue[SourceShape[T], (Sink[T, NotUsed],
MergeHub.DrainingControl)] {
require(perProducerBufferSize > 0, "Buffer size must be positive")
+ require(maxTotalBufferSize >= 0, "Max total buffer size must be non-negative
(0 means unlimited)")
val out: Outlet[T] = Outlet("MergeHub.out")
override val shape: SourceShape[T] = SourceShape(out)
@@ -182,6 +232,7 @@ private[pekko] class MergeHub[T](perProducerBufferSize:
Int, drainingEnabled: Bo
* processing of control messages. This causes no issues though, see the
explanation in 'tryProcessNext'.
*/
private val queue = new AbstractNodeQueue[Event] {}
+ private val totalBufferedElements = new AtomicInteger(0)
@volatile private var needWakeup = false
@volatile private var shuttingDown = false
@volatile private var draining = false
@@ -224,6 +275,7 @@ private[pekko] class MergeHub[T](perProducerBufferSize:
Int, drainingEnabled: Bo
needWakeup = false
nextElem match {
case Element(id, elem) =>
+ totalBufferedElements.decrementAndGet()
demands(id).onElement()
push(out, elem)
// demand consumed by push — exit and wait for the next onPull
@@ -248,9 +300,14 @@ private[pekko] class MergeHub[T](perProducerBufferSize:
Int, drainingEnabled: Bo
def isShuttingDown: Boolean = shuttingDown
def isDraining: Boolean = drainingEnabled && draining
+ def isBufferFull: Boolean = maxTotalBufferSize > 0 &&
totalBufferedElements.get() >= maxTotalBufferSize
// External API
private[MergeHub] def enqueue(ev: Event): Unit = {
+ ev match {
+ case _: Element => totalBufferedElements.incrementAndGet()
+ case _ =>
+ }
queue.add(ev)
/*
* Simple volatile var is enough, there is no need for a CAS here. The
first important thing to note
@@ -325,7 +382,8 @@ private[pekko] class MergeHub[T](perProducerBufferSize:
Int, drainingEnabled: Bo
private val id = idCounter.getAndIncrement()
override def preStart(): Unit = {
- if (!logic.isDraining && !logic.isShuttingDown) {
+ // Best-effort admission check: already admitted producers may
overshoot the bound.
+ if (!logic.isDraining && !logic.isShuttingDown &&
!logic.isBufferFull) {
logic.enqueue(Register(id, getAsyncCallback(onDemand)))
// At this point, we could be in the unfortunate situation that:
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]