为何未触发Failure分支?如何捕获Akka的StreamTcpException?
Let’s break down exactly what’s happening here and how to fix it:
The Root Cause
Your RestartSource.withBackoff is doing exactly what it’s designed to do: it swallows internal stream failures and automatically restarts the stream according to your backoff policy.
When the WebSocket connection fails (like when the server is down), the internal stream throws a StreamTcpException. But instead of letting this failure propagate to the outer Future returned by .runWith(), RestartSource catches it, schedules a restart, and keeps the stream lifecycle going. That’s why your onComplete’s Failure branch never fires— the outer Future only completes successfully if the stream is gracefully terminated (e.g., when you shut down the ActorSystem) or fails only if something catastrophic breaks the RestartSource itself (not regular connection errors).
The StreamTcpException warning you see is just Akka Streams logging the internal failure, but RestartSource is handling it to trigger the restart logic.
How to Capture Connection Failures
Instead of relying on the outer Future’s onComplete, you need to handle failures inside the RestartSource’s internal stream. Here are two clean, practical approaches:
Approach 1: Handle Failures Directly in the Internal Source
Add a recover stage to catch stream failures and log your custom message before re-throwing the exception to trigger the restart:
object Main extends App { private val sapServer = "127.0.0.1:8080" implicit val system = ActorSystem("WsSystem") implicit val materializer = ActorMaterializer() implicit val dispatcher = system.dispatcher RestartSource.withBackoff( minBackoff = 3.seconds, maxBackoff = 30.seconds, randomFactor = 0.2 ) { () => val (supported, source) = Source.tick(1.seconds, 15.seconds, TextMessage.Strict("")) .viaMat(Http().webSocketClientFlow(WebSocketRequest(s"ws://$sapServer")))(Keep.right) .preMaterialize() // Log upgrade failures immediately supported.onFailure { case ex => println("Probably server is down.") } // Catch and log stream-level failures like StreamTcpException source.recover { case ex: Exception => println(s"Connection failed: ${ex.getMessage}") throw ex // Re-throw to trigger RestartSource's restart logic } }.runWith(Sink.foreach(println)) }
Approach 2: Use a Sink with Custom Supervision
Modify your sink to include a supervision strategy that logs errors and restarts the stream:
object Main extends App { private val sapServer = "127.0.0.1:8080" implicit val system = ActorSystem("WsSystem") implicit val materializer = ActorMaterializer() implicit val dispatcher = system.dispatcher // Create a sink with error handling logic val errorAwareSink = Sink.foreach[Message](println) .withAttributes(ActorAttributes.supervisionStrategy { case _: Exception => println("Probably server is down.") Supervision.Restart // Trigger stream restart via RestartSource case _ => Supervision.Stop }) RestartSource.withBackoff( minBackoff = 3.seconds, maxBackoff = 30.seconds, randomFactor = 0.2 ) { () => val (supported, source) = Source.tick(1.seconds, 15.seconds, TextMessage.Strict("")) .viaMat(Http().webSocketClientFlow(WebSocketRequest(s"ws://$sapServer")))(Keep.right) .preMaterialize() supported.flatMap { upgrade => if (upgrade.response.status == StatusCodes.SwitchingProtocols) Future.successful(Done) else throw new RuntimeException(s"Connection failed: ${upgrade.response.status}") } source }.runWith(errorAwareSink) }
Key Notes
RestartSourceis built to keep streams alive by restarting them on failure, so the outerFuturewill almost never enter theFailurestate. All stream-specific errors should be handled within the stream’s lifecycle.- If you want to stop restarting after a certain number of failures, you can track restart counts manually (e.g., using a mutable variable or stateful stream stage) and adjust the backoff logic to terminate the stream once the limit is hit.
内容的提问来源于stack exchange,提问作者softshipper

