This is an automated email from the ASF dual-hosted git repository.

He-Pin pushed a commit to branch feat/mergehub-max-total-buffer-size
in repository https://gitbox.apache.org/repos/asf/pekko.git

commit a40eaff369800b55c6e7e56d6c8a95a612b7b627
Author: θ™ŽιΈ£ <[email protected]>
AuthorDate: Mon Jul 27 23:58:08 2026 +0800

    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
---
 .../org/apache/pekko/stream/scaladsl/HubSpec.scala | 60 +++++++++++++++++++++
 .../org/apache/pekko/stream/javadsl/Hub.scala      | 62 +++++++++++++++++++++
 .../org/apache/pekko/stream/scaladsl/Hub.scala     | 63 ++++++++++++++++++++--
 3 files changed, 181 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..68782daa21 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,66 @@ class HubSpec extends StreamSpec {
       sub.cancel()
     }
 
+    "cancel new producers when maxTotalBufferSize is exceeded" in {
+      val downstream = TestSubscriber.manualProbe[Int]()
+      // perProducerBufferSize=4, maxTotalBufferSize=8
+      val sink = 
Sink.fromSubscriber(downstream).runWith(MergeHub.source[Int](4, 8))
+
+      val sub = downstream.expectSubscription()
+      // Do not request any elements β€” downstream is stalled
+
+      // First producer connects and pushes up to 4 elements
+      val upstream1 = TestPublisher.probe[Int]()
+      Source.fromPublisher(upstream1).runWith(sink)
+      for (i <- 1 to 4) upstream1.sendNext(i)
+
+      // Second producer connects and pushes up to 4 elements (total = 8 = 
maxTotalBufferSize)
+      val upstream2 = TestPublisher.probe[Int]()
+      Source.fromPublisher(upstream2).runWith(sink)
+      for (i <- 5 to 8) upstream2.sendNext(i)
+
+      // Third producer should be cancelled because buffer is full
+      val upstream3 = TestPublisher.probe[Int]()
+      Source.fromPublisher(upstream3).runWith(sink)
+      upstream3.expectCancellation()
+
+      sub.cancel()
+    }
+
+    "allow new producers after elements are consumed from the aggregate 
buffer" in {
+      val downstream = TestSubscriber.manualProbe[Int]()
+      // perProducerBufferSize=4, maxTotalBufferSize=4
+      val sink = 
Sink.fromSubscriber(downstream).runWith(MergeHub.source[Int](4, 4))
+
+      val sub = downstream.expectSubscription()
+
+      // First producer connects and pushes 4 elements (fills the buffer)
+      val upstream1 = TestPublisher.probe[Int]()
+      Source.fromPublisher(upstream1).runWith(sink)
+      for (i <- 1 to 4) upstream1.sendNext(i)
+
+      // Second producer is cancelled because buffer is full
+      val upstream2 = TestPublisher.probe[Int]()
+      Source.fromPublisher(upstream2).runWith(sink)
+      upstream2.expectCancellation()
+
+      // Now consume all 4 elements, freeing the buffer
+      sub.request(4)
+      downstream.expectNext(1)
+      downstream.expectNext(2)
+      downstream.expectNext(3)
+      downstream.expectNext(4)
+
+      // Now a new producer should be accepted
+      val upstream3 = TestPublisher.probe[Int]()
+      Source.fromPublisher(upstream3).runWith(sink)
+      upstream3.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..2a377e4a49 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,34 @@ 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.
+   *   Transient overshoot up to the per-producer buffer of concurrently 
admitted producers is possible.
+   *   Use 0 for unlimited (default behavior).
+   */
+  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 +126,40 @@ 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 produces 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. Default 
value is 16.
+   * @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.
+   *   Transient overshoot up to the per-producer buffer of concurrently 
admitted producers is possible.
+   *   Use 0 for unlimited (default behavior).
+   */
+  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..1f0a3cb84a 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,26 @@ 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. Default 
value is 16.
+   * @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.
+   *   Transient overshoot up to the per-producer buffer of concurrently 
admitted producers is possible.
+   *   Use 0 for unlimited (default behavior).
+   */
+  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 +116,31 @@ 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. Default 
value is 16.
+   * @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.
+   *   Transient overshoot up to the per-producer buffer of concurrently 
admitted producers is possible.
+   *   Use 0 for unlimited (default behavior).
+   */
+  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 +184,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 +230,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 +273,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 +298,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 +380,7 @@ private[pekko] class MergeHub[T](perProducerBufferSize: 
Int, drainingEnabled: Bo
           private val id = idCounter.getAndIncrement()
 
           override def preStart(): Unit = {
-            if (!logic.isDraining && !logic.isShuttingDown) {
+            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