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]
