Akka Stream, Tcp().bind, 客户端关闭socket时的处理

问题描述

我是 Akka Stream 的新手,我想了解如何为我的项目处理 TCP 套接字。我从 Akka Stream official documentation获取了这段代码

import akka.stream.scaladsl.Framing

val connections: Source[IncomingConnection,Future[ServerBinding]] =
  Tcp().bind(host,port)

connections.runForeach { connection =>
  println(s"New connection from: ${connection.remoteAddress}")

  val echo = Flow[ByteString]
    .via(Framing.delimiter(ByteString("\n"),maximumFrameLength = 256,allowTruncation = true))
    .map(_.utf8String)
    .map(_ + "!!!\n")
    .map(ByteString(_))

  connection.handleWith(echo)
}

如果我使用 netcat 从终端连接,我可以看到 Akka Stream TCP 套接字按预期工作。我还发现如果我需要使用用户消息关闭连接,我可以使用 takeWhile 如下

import akka.stream.scaladsl.Framing

val connections: Source[IncomingConnection,allowTruncation = true))
    .map(_.utf8String)
    .takeWhile(_.toLowerCase.trim != "exit")   // < - - - - - - HERE
    .map(_ + "!!!\n")
    .map(ByteString(_))

  connection.handleWith(echo)
}

我找不到的是如何管理由 CMD + C 操作关闭套接字。 Akka Stream 在内部使用 Akka.io 来管理 TCP 连接,因此它必须在套接关闭时发送一些 PeerClose 消息。所以,我对 Akka.io 的理解告诉我,我应该收到来自套接关闭的反馈,但我找不到如何使用 Akka Stream 做到这一点。有没有办法管理?

解决方法

connection.handleWith(echo)connection.flow.joinMat(echo)(Keep.right).run() 的语法糖,它具有 echo 的物化值,通常没有用。 Flow.via.map.takeWhileNotUsed 作为物化值,所以这也基本没用。但是,您可以将阶段附加到 echo,这会以不同的方式实现。

其中之一是 .watchTermination

connections.runForeach { connection =>
  println(s"New connection from: ${connection.remoteAddress}")

  val echo: Flow[ByteString,ByteString,Future[Done]] = Flow[ByteString]
    .via(Framing.delimiter(ByteString("\n"),maximumFrameLength = 256,allowTruncation = true))
    .map(_.utf8String)
    .takeWhile(_.toLowerCase.trim != "exit")   // < - - - - - - HERE
    .map(_ + "!!!\n")
    .map(ByteString(_))
    // change the materialized value to a Future[Done]
    .watchTermination()(Keep.right)

  // you may need to have an implicit ExecutionContext in scope,e.g. system.dispatcher,//  if you don't already
  connection.handleWith(echo).onComplete {
    case Success(_) => println("stream completed successfully")
    case Failure(e) => println(e.getMessage)
  }
}

这不会区分你这边还是远端正常关闭连接;它将区分流失败。

相关问答

Selenium Web驱动程序和Java。元素在(x,y)点处不可单击。其...
Python-如何使用点“。” 访问字典成员?
Java 字符串是不可变的。到底是什么意思?
Java中的“ final”关键字如何工作?(我仍然可以修改对象。...
“loop:”在Java代码中。这是什么,为什么要编译?
java.lang.ClassNotFoundException:sun.jdbc.odbc.JdbcOdbc...