【问题标题】:How can I read several web-sockets within a single Akka-http Flow?如何在单个 Akka-http 流中读取多个 Web 套接字?
【发布时间】:2020-07-13 06:52:56
【问题描述】:

我目前正在通过尝试建立多个 websockets 连接来练习 Akka-http。我创建 websockets 客户端流 (sn-p) 的代码如下所示:

val webSocketFlow =
  Http().webSocketClientFlow(WebSocketRequest(url), settings = customSettings)

val (upgradeResponse, closed) =
  outgoing
    .viaMat(webSocketFlow)(Keep.right)
    .viaMat(decoder)(Keep.left)
    .toMat(sink)(Keep.both)
    .run()

如果我有一个网址,这目前效果很好。我很好奇如何扩展它以连接到多个 url。例如,如果我有一个不确定的 websockets 端点列表List("ws://localhost:8080/foo", "ws://localhost:8080/bar", "ws://localhost:8080/baz")

我已经考虑为每个 URL 添加一个新流,但是如果我有很长的 websockets 端点/url 列表怎么办。然后,这变得繁琐且过于手动。我还考虑过将其包装到一个函数中,并在给定的迭代中调用每个 URL。但这也让人觉得有点过头了。

有没有办法让一个连接池全部来源于一个流(或类似的东西)?也欢迎进一步阅读。作为“不错的选择”,是否还有一种方法可以标记传入的消息以使用它们来自的 url 发出信号?

更新:为澄清起见,我只是从 websockets 读取(仅客户端),而不是发回任何消息。

【问题讨论】:

  • 您是向每个 WebSocket 发送相同的内容,还是向每个 WebSocket 发送不同的消息?
  • @Tim 我已经更新了我的答案以进行澄清。我只从 websockets 阅读。所以想法是从多个 websocket 端点读取

标签: scala websocket akka akka-http


【解决方案1】:

这应该可以工作(代码正在文本框中编写...):

def taggedWebsocketForUrl(url: String, tag: Int): Source[(Int, Message), Future[WebSocketUpgradeResponse]] =
  outgoing.viaMat(Http().webSocketClientFlow(WebSocketRequest(url), settings = customSettings))(Keep.right).map(tag -> _)

val websocketMergedSource: Source[(Int, Message), Seq[Future[WebSocketUpgradeResponse]]] = {
  // You could replace this with a mess of headOptions etc., but...
  if (websocketUrls.isEmpty) Source.empty[(Int, Message)].mapMaterializedValue(_ => Seq(Future.failed(new NoSuchElementException("no websocket URLs"))))
  else {
    val first: Source[(Int, Message), List[Future[WebSocketUpgradeResponse]]] =
      taggedWebsocketForUrl(websocketUrls.head, 0).mapMaterializedValue(List(_))
    if (websocketUrls.tail.isEmpty) first
    else {
      websocketUrls.tail.foldLeft(first -> 1) {
        (acc, url) =>
          val newSource = acc._1.mergeMat(taggedWebsocketForUrl(url, acc._2)) {
            (futs: List[Future[WebSocketUpgradeResponse]], fut: Future[WebSocketUpgradeResponse]) =>
              fut :: futs // Will reverse at the end...
          }
          newSource -> (acc._2 + 1)
      }._1.mapMaterializedValue(_.reverse)
    }
  }
}

有了这个,您将有许多升级响应(您可以mapMaterializedValue(Future.sequence _) 将它们组合成一个Future[Seq[WebsocketUpgradeResponse]],如果任何失败都会失败)。来自列表中nth url 的消息将被标记为n

请注意,websocketUrlsList 引导构建为折叠:如果有 n url,来自第一个 url 的消息将经过 n-1 合并阶段,最后一个 url 将仅经历 1 个合并阶段,因此您希望将期望产生更多流量的网址放在列表末尾。

另一种更有效的方法是使用IndexedSeq(如VectorArray)来分而治之,以建立merges 的树。

使用 Akka Streams GraphDSL 也会给你很多控制权,但我倾向于仅将其用作最后的手段。

【讨论】:

  • Point Levi 的另一个答案!谢谢。我对您的回答做了一些小修改,如果它们偏离了您的意图,请告诉我。我设法使代码工作,但在mergeMat(taggedWebsocketForUrl(url, acc._2)) 之后失去了你。据我了解,我们将所有来源合并在一起,然后函数 curry 的第二部分让我有些困惑。 `fut :: futs`的意图是什么,为什么我们需要在最后反转?我也很高兴在一个单独的问题中讨论这个问题,如果它对其他人来说更清晰的参考,mergeMat 是如何工作的
  • mergeMat 就像merge (将两个输入组合成一个输出;如果两个输入都有可用的元素,则随机选择)但merge 保留物化值(在这种情况下,左侧的Future[WebsocketUpgradeResponse]]),mergeMat 为我们提供了两个物化值,我们提供了一个函数来创建一个新的物化值(在这种情况下,将所有物化值累积在一个 Seq 中)。跨度>
  • 至于最后的fut :: futsreverse,这是因为通过前置构建列表比通过附加构建列表要快得多,尤其是在元素很多的情况下。前置的缺点是结果列表的顺序是相反的,所以最后的 reverse 会使列表以正确的顺序排列。
【解决方案2】:

您需要某种方法来合并各种 WebSocket 流,以便您可以处理传入的消息,就好像它们来自单个源一样。

由于您不需要发送任何数据而只需要接收实现就很简单了。

让我们开始创建一个函数,该函数将为给定的 uri 创建一个 WebSocket 源:

def webSocketSource(uri: Uri): Source[Message, Future[WebSocketUpgradeResponse]] = {
  Source.empty.viaMat(Http().webSocketClientFlow(uri))(Keep.right)
}

由于您不关心发送数据,该函数通过提供一个空的 Source 立即关闭输出通道。结果是一个 Source,其中包含从 WebSocket 读取的消息。

此时我们可以使用这个函数为每个uri创建一个专用的源:

val wsSources: List[Source[Message, NotUsed]] = uris.map { uri =>
  webSocketSource(uri).mapMaterializedValue { respFuture =>
    respFuture.map {
      case _: ValidUpgrade => log.debug(s"Websocket upgrade for [${uri}] successful")
      case err: InvalidUpgradeResponse => log.error(s"Websocket upgrade for [${uri}] failed: ${err.cause}")
    }

    NotUsed
  }
}

在这里,我们需要以某种方式处理具体化的值,因为我们不知道它们有多少,所以不可能(或至少不容易)组合它们。因此,我们在这里采用最简单的方法,即只记录日志。

现在我们已经准备好我们的源,我们可以继续合并它们:

val mergedSource: Source[Message, NotUsed] = wsSources match {
  case s1 :: s2 :: rest => Source.combine(s1, s2, rest: _*)(Merge(_))
  case s1 :: Nil => s1
  case Nil => Source.empty[Message]
}

这里的想法是,如果我们有 2 个或更多 uri,我们实际上会执行合并操作,否则如果我们只有一个,我们只需使用它而不做任何修改。最后,我们还通过提供一个空的 Source 来涵盖根本没有任何 uri 的情况,该 Source 将简单地终止流而不会出错。

此时我们可以将此源与您已有的流和接收器结合起来并运行它:

val done: Future[Done] = mergedSource.via(decoder).toMat(sink)(Keep.right).run

这给了我们一个单一的未来,它将在所有连接完成时完成或在一个连接失败时立即失败。

【讨论】:

    猜你喜欢
    • 2016-10-14
    • 2012-05-16
    • 1970-01-01
    • 1970-01-01
    • 2020-10-20
    • 1970-01-01
    • 2018-05-29
    • 2011-05-31
    • 2021-12-08
    相关资源
    最近更新 更多