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]

Reply via email to