【问题标题】:How to recover from akka.stream.io.Framing$FramingException如何从 akka.stream.io.Framing$FramingException 中恢复
【发布时间】:2015-07-23 19:44:36
【问题描述】:

开启:akka-stream-experimental_2.11 1.0。

我们在 Tcp 服务器中使用 Framing.delimiter。当消息到达的长度大于 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 不尊重supervision

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

【问题讨论】:

    标签: scala-2.11 akka-stream


    【解决方案1】:

    您是否尝试过对 targetFlow 接收器而不是整个物化器进行监督?我在这里的任何地方都看不到它,我认为它应该直接设置在该流上。

    这更像是一个猜测而不是 科学 ;)

    【讨论】:

      【解决方案2】:

      我从文件中读取时遇到了同样的异常,对我来说,这是通过在最后一行后面加上一个 return 来解决的。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2022-11-23
        • 2012-09-16
        • 2019-06-07
        • 2012-10-06
        • 2017-01-22
        相关资源
        最近更新 更多