You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何未触发Failure分支?如何捕获Akka的StreamTcpException?

Why the Failure Branch Isn't Triggered & How to Catch the Exception

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

  • RestartSource is built to keep streams alive by restarting them on failure, so the outer Future will almost never enter the Failure state. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 08:14:13