This is an automated email from the ASF dual-hosted git repository. He-Pin pushed a commit to branch fix/broadcast-hub-poststop in repository https://gitbox.apache.org/repos/asf/pekko.git
commit c751b6f102d1ff00e2a2010133f29d0090eab58a Author: θιΈ£ <[email protected]> AuthorDate: Tue Jul 14 18:45:44 2026 +0800 fix(stream): notify BroadcastHub consumers on materializer shutdown Motivation: When the materializer running a BroadcastHub.sink is shut down (e.g. the hosting actor stops), consumers materialized on a different materializer may not be notified of the hub's termination (see #3345). Modification: Move registered-consumer (consumerWheel) notification from onUpstreamFailure to postStop so postStop is the single notification point. postStop now handles both Open (normal termination) and Closed (after upstream failure) states, ensuring registered consumers always receive the appropriate signal. Result: BroadcastHub consumers are reliably notified when the hub's materializer shuts down, preventing consumers from hanging indefinitely. Tests: - sbt "stream-tests / Test / testOnly org.apache.pekko.stream.scaladsl.HubSpec" 51 passed, 0 failed References: Fixes #3345 --- .../org/apache/pekko/stream/scaladsl/HubSpec.scala | 24 +++++++++++++- .../org/apache/pekko/stream/scaladsl/Hub.scala | 37 ++++++++++++++-------- 2 files changed, 47 insertions(+), 14 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 b32c8746b1..628585cca0 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 @@ -19,7 +19,7 @@ import scala.concurrent.duration._ import org.apache.pekko import pekko.Done -import pekko.stream.KillSwitches +import pekko.stream.{ KillSwitches, Materializer } import pekko.stream.ThrottleMode import pekko.stream.impl.ActorPublisher import pekko.stream.testkit.StreamSpec @@ -559,6 +559,28 @@ class HubSpec extends StreamSpec { } } + "notify consumers when hub materializer is shut down" in { + val hubMat = Materializer(system) + val upstream = TestPublisher.probe[Int]() + val source = Source.fromPublisher(upstream).runWith(BroadcastHub.sink(8))(hubMat) + + val downstream = TestSubscriber.probe[Int]() + source.runWith(Sink.fromSubscriber(downstream)) + + downstream.request(1) + Thread.sleep(100) + + upstream.sendNext(1) + downstream.expectNext(1) + + hubMat.shutdown() + + // The consumer must be notified (not left hanging), which is the fix for #3345. + // The signal is an error because the SubSink output boundary failure + // arrives before the postStop completion callback. + downstream.expectError() + } + "handle cancelled Sink" in { val in = TestPublisher.probe[Int]() val hubSource = Source.fromPublisher(in).runWith(BroadcastHub.sink(4)) 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 32e647841f..f4ce5cdebd 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 @@ -652,21 +652,12 @@ private[pekko] class BroadcastHub[T](startAfterNrOfConsumers: Int, bufferSize: I override def onUpstreamFailure(ex: Throwable): Unit = { val failMessage = HubCompleted(Some(ex)) - // Notify pending consumers and set tombstone + // Notify pending consumers and set tombstone. + // Registered consumers in the consumerWheel are notified by postStop. state.getAndSet(Closed(Some(ex))).asInstanceOf[Open].registrations.foreach { consumer => consumer.callback.invoke(failMessage) } - // Notify registered consumers β skip null (empty) slots - var idx = 0 - while (idx < consumerWheel.length) { - val bucket = consumerWheel(idx) - if (bucket ne null) { - val itr = bucket.valuesIterator - while (itr.hasNext) itr.next().callback.invoke(failMessage) - } - idx += 1 - } failStage(ex) } @@ -759,19 +750,39 @@ private[pekko] class BroadcastHub[T](startAfterNrOfConsumers: Int, bufferSize: I } override def postStop(): Unit = { - // Notify pending consumers and set tombstone + // Notify all consumers (pending and registered) when the stage stops. + // Registered consumers in the consumerWheel are notified here (not in onUpstreamFailure) + // so that materializer shutdown produces the correct signal. @tailrec def tryClose(): Unit = state.get() match { - case Closed(_) => // Already closed, ignore + case Closed(_) => // Already closed by onUpstreamFailure β fall through to notify wheel below + notifyRegisteredConsumers() case open: Open => if (state.compareAndSet(open, Closed(None))) { val completedMessage = HubCompleted(None) open.registrations.foreach { consumer => consumer.callback.invoke(completedMessage) } + notifyRegisteredConsumers() } else tryClose() } + def notifyRegisteredConsumers(): Unit = { + val message = state.get() match { + case Closed(Some(ex)) => HubCompleted(Some(ex)) + case _ => HubCompleted(None) + } + var idx = 0 + while (idx < consumerWheel.length) { + val bucket = consumerWheel(idx) + if (bucket ne null) { + val itr = bucket.valuesIterator + while (itr.hasNext) itr.next().callback.invoke(message) + } + idx += 1 + } + } + tryClose() } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
