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 b3c7c24fcc fix(stream): notify PartitionHub consumers on materializer 
shutdown (#3347)
b3c7c24fcc is described below

commit b3c7c24fccb386f81930f1a25b27cc866cd1dc92
Author: He-Pin(kerr) <[email protected]>
AuthorDate: Wed Jul 15 20:49:09 2026 +0800

    fix(stream): notify PartitionHub consumers on materializer shutdown (#3347)
    
    Motivation:
    When the materializer running a PartitionHub.sink is shut down (e.g. the
    hosting actor stops or crashes), consumers materialized on a different
    materializer may not be notified of the hub's termination. This mirrors
    the same bug pattern found in BroadcastHub (see #3345).
    
    Modification:
    Move registered-consumer notification from onUpstreamFailure to postStop
    so that postStop is the single notification point for all consumers.
    postStop now handles both Open (normal termination) and Closed (after
    upstream failure) states, ensuring registered consumers always receive
    the appropriate signal. Each callback invocation is wrapped in
    try/catch NonFatal so a single consumer failure does not short-circuit
    the notification loop.
    
    Result:
    PartitionHub 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"
    
    References:
    Refs #3345
---
 .../org/apache/pekko/stream/scaladsl/HubSpec.scala | 23 ++++++++++++++++
 .../org/apache/pekko/stream/scaladsl/Hub.scala     | 32 ++++++++++++++++------
 2 files changed, 46 insertions(+), 9 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 954c592da9..deaebb4d51 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
@@ -968,6 +968,29 @@ class HubSpec extends StreamSpec {
       result2.futureValue should ===((1 to 9 by 2).filterNot(_ == 3))
     }
 
+    "notify consumers when hub materializer is shut down" in {
+      val hubMat = Materializer(system)
+      val upstream = TestPublisher.probe[Int]()
+      val source = Source.fromPublisher(upstream).runWith(
+        PartitionHub.sink((size, elem) => elem % size, startAfterNrOfConsumers 
= 1, bufferSize = 8))(hubMat)
+
+      val downstream = TestSubscriber.probe[Int]()
+      source.runWith(Sink.fromSubscriber(downstream))
+
+      downstream.request(1)
+      upstream.expectRequest()
+
+      upstream.sendNext(0)
+      downstream.expectNext(0)
+
+      hubMat.shutdown()
+
+      // The consumer must be notified (not left hanging), which is the fix.
+      // The signal is an error because the SubSink output boundary failure
+      // arrives before the postStop completion callback.
+      downstream.expectError()
+    }
+
   }
 
 }
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 ca7615a55f..26720b05ec 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
@@ -1374,32 +1374,46 @@ object PartitionHub {
     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 are notified by postStop.
       
state.getAndSet(Closed(Some(ex))).asInstanceOf[Open].registrations.foreach { 
consumer =>
-        consumer.callback.invoke(failMessage)
+        try consumer.callback.invoke(failMessage)
+        catch { case NonFatal(_) => }
       }
 
-      // Notify registered consumers
-      consumerInfo.consumers.foreach { consumer =>
-        consumer.callback.invoke(failMessage)
-      }
       failStage(ex)
     }
 
     override def postStop(): Unit = {
-      // Notify pending consumers and set tombstone
+      // Notify all consumers (pending and registered) when the stage stops.
+      // Registered consumers 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 — notify 
registered consumers here
+          notifyRegisteredConsumers()
         case open: Open =>
           if (state.compareAndSet(open, Closed(None))) {
             val completedMessage = HubCompleted(None)
             open.registrations.foreach { consumer =>
-              consumer.callback.invoke(completedMessage)
+              try consumer.callback.invoke(completedMessage)
+              catch { case NonFatal(_) => }
             }
+            notifyRegisteredConsumers()
           } else tryClose()
       }
 
+      def notifyRegisteredConsumers(): Unit = {
+        val message = state.get() match {
+          case Closed(Some(ex)) => HubCompleted(Some(ex))
+          case _                => HubCompleted(None)
+        }
+        consumerInfo.consumers.foreach { consumer =>
+          try consumer.callback.invoke(message)
+          catch { case NonFatal(_) => }
+        }
+      }
+
       tryClose()
     }
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to