如何从akka.stream.io.Framing$FramingException中恢复



On: akka-stream-experimental_2.11 1.0.

我们在Tcp服务器中使用帧分隔符。当到达的消息长度大于maximumFrameLength时,会抛出FramingException,我们可以从ActorSubscriber的OnError中捕获它。

服务器代码:

def bind(address: String, port: Int, target: ActorRef, maxInFlight: Int, maxFrameLength: Int)
    (implicit system: ActorSystem, actorMaterializer: ActorMaterializer): Future[ServerBinding] = {
    val sink = Sink.foreach {
      conn: Tcp.IncomingConnection =>
        val targetSubscriber = ActorSubscriber[Message](system.actorOf(Props(new TargetSubscriber(target, maxInFlight))))
        val targetSink = Flow[ByteString]
          .via(Framing.delimiter(ByteString("n"), maximumFrameLength = maxFrameLength, allowTruncation = true))
          .map(raw ⇒ Message(raw))
          .to(Sink(targetSubscriber))
        conn.flow.to(targetSink).runWith(Source(Promise().future))
    }
    val connections = Tcp().bind(address, port)
    connections.to(sink).run()
  }
用户代码:

class TargetSubscriber(target: ActorRef, maxInFlight: Int) extends ActorSubscriber with ActorLogging {
  private var inFlight = 0
  override protected def requestStrategy = new MaxInFlightRequestStrategy(maxInFlight) {
    override def inFlightInternally = inFlight
  }
  override def receive = {
    case OnNext(msg: Message) ⇒
      target ! msg
      inFlight += 1
    case OnError(t) ⇒
      inFlight -= 1
      log.error(t, "Subscriber encountered error")
    case TargetAck(_) ⇒
      inFlight -= 1
  }
}

问题:在该传入连接的此异常之后,低于最大帧长度的消息不会流。终止客户端并重新运行它可以正常工作。

ActorSubscriber不遵守监督

跳过坏消息,继续下一个好消息的正确方法是什么?

您是否尝试对targetFlow sink而不是整个materialiser进行监督?我在这里没有看到它,我认为它应该直接设置在那个流上。

比起science;)

我从一个文件中读取同样的异常,对我来说,通过在最后一行后面放一个返回来解决这个问题。

相关内容

  • 没有找到相关文章

最新更新