【问题标题】:Closing an Akka stream from inside a GraphStage (Akka 2.4.2)从 GraphStage (Akka 2.4.2) 内部关闭 Akka 流
【发布时间】:2016-02-24 11:54:41
【问题描述】:

在 Akka Stream 2.4.2 中,PushStage 已被弃用。对于 Streams 2.0.3,我使用了这个答案中的解决方案:

How does one close an Akka stream?

原来是:

import akka.stream.stage._

    val closeStage = new PushStage[Tpe, Tpe] {
      override def onPush(elem: Tpe, ctx: Context[Tpe]) = elem match {
        case elem if shouldCloseStream ⇒
          // println("stream closed")
          ctx.finish()
        case elem ⇒
          ctx.push(elem)
      }
    }

如何从 GraphStage / onPush() 内部立即关闭 2.4.2 中的流?

【问题讨论】:

    标签: scala akka-stream


    【解决方案1】:

    使用这样的东西:

    val closeStage = new GraphStage[FlowShape[Tpe, Tpe]] {
      val in = Inlet[Tpe]("closeStage.in")
      val out = Outlet[Tpe]("closeStage.out")
    
      override val shape = FlowShape.of(in, out)
    
      override def createLogic(inheritedAttributes: Attributes) = new GraphStageLogic(shape) {
        setHandler(in, new InHandler {
          override def onPush() = grab(in) match {
            case elem if shouldCloseStream ⇒
              // println("stream closed")
              completeStage()
            case msg ⇒
              push(out, msg)
          }
        })
        setHandler(out, new OutHandler {
          override def onPull() = pull(in)
        })
      }
    }
    

    它更冗长,但一方面可以以可重用的方式定义此逻辑,另一方面不必担心流元素之间的差异,因为GraphStage 可以以相同的方式处理作为一个流程将被处理:

    val flow: Flow[Tpe] = ???
    val newFlow = flow.via(closeStage)
    

    【讨论】:

    • 这就是我最初尝试的方法,但它似乎无法正常工作。在我执行 completeStage() 并关闭客户端后,一分钟后我的服务器上出现异常,“由于上游故障而中止 TCP 连接:最后 1 分钟内没有传递任何元素。” pastebin.com/utu11L2h
    • 当然,completeStage 只关闭客户端上的流,而不是服务器上的流。在调用completeStage 之前,您可以向服务器发送一条消息,告知您将关闭客户端并且它可以关闭连接。
    • 我认为这就是我所缺少的,我没有看到大局,如果这听起来很愚蠢,对不起。当您说“向服务器发送消息”时,它是如何工作的?它不可能是来自客户端说“我要离开”的任何东西,因为我正在关闭连接,因为客户端行为不端(锁定、滞后、错误实施、协议违规等)并且可能不会发送“我”米离开”。 Close 需要由 GraphStage 驱动服务器。
    • 好吧,如果你失去了与服务器的连接,那么你在客户端就无能为力了。服务器上的故障是什么问题?客户走了,失败告诉你。
    • 我认为我的问题是a)我不喜欢未处理的异常或一般的异常,并且b)仍然有资源使用了1分钟。但关键是,客户还没走,不管它要不要合作,我都想把它踢掉。例如,滞后或挂起的客户端仍可能有连接,只是无法发送或接收数据。
    【解决方案2】:

    发帖供他人参考。 sschaef 的回答在程序上是正确的,但连接保持打开一分钟,最终会超时并抛出“无活动”异常,关闭连接。

    在进一步阅读文档时,我注意到当所有上游流程完成时连接已关闭。就我而言,我有不止一个上游。

    对于我的特定用例,解决方法是添加 eagerComplete=true 以在任何(而不是全部)上游完成后立即关闭流。比如:

    ... = builder.add(Merge[MyObj](3,eagerComplete = true))
    

    希望这对某人有所帮助。

    【讨论】:

    • 感谢您的参考。我还建议将@sschaef 的答案标记为正确,因为它可以帮助那些有“如何关闭流程”问题而没有“eagerComplete”问题的人(比如我)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-04-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多