You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从Akka Http WebSocket客户端接收消息并按需回复?

看起来你是在用Akka Streams(大概率配合Akka HTTP WebSocket)实现客户端消息推送,现在要扩展成双向通信——接收客户端发来的消息并给出回复对吧?我来给你拆解下怎么改:

首先,你原来的代码是单向的 outgoing 流,而WebSocket需要的是Flow[Message, Message, _],这个Flow要同时处理两个方向:

  • 上游(incoming):客户端发给我们的消息
  • 下游(outgoing):我们发给客户端的消息(包括主动推送和回复)

第一步:处理客户端的 incoming 消息

先写一个函数,把客户端的消息转换成我们的回复。要注意Akka的Message分几种类型,别漏了:

import akka.NotUsed
import akka.stream.scaladsl.{Flow, Source, Sink}
import akka.http.scaladsl.model.ws.{Message, TextMessage, BinaryMessage}
import scala.concurrent.Future

def processClientMessage(msg: Message): Future[Message] = msg match {
  // 处理一次性发来的文本消息
  case TextMessage.Strict(clientText) =>
    // 这里写你的业务逻辑:比如解析clientText,生成对应回复
    Future.successful(TextMessage(s"Got your message: '$clientText' — thanks!"))
  
  // 处理流式文本消息(客户端分块发的)
  case TextMessage.Streamed(textStream) =>
    // 先把流收集成完整文本,再处理
    textStream.runFold("")(_ + _).map(fullText => 
      TextMessage(s"Received streamed content (length: ${fullText.length}): $fullText")
    )
  
  // 处理二进制消息(如果不需要可以返回提示)
  case BinaryMessage(_) =>
    Future.successful(TextMessage("Sorry, I don't handle binary messages right now!"))
}

第二步:构建双向Flow

现在把你原来的主动推送逻辑和回复逻辑结合起来:

  • 原来的source不要用runWith(Sink.ignore),那会把数据扔掉,直接保留成Source[Message, _]
  • 用mapAsync处理客户端消息生成回复,再和主动推送的流合并

完整代码示例:

// 你的主动推送数据源(保留原来的throttle逻辑)
val activePushSource: Source[Message, NotUsed] = 
  Source.fromFuture(data) // 这里的data应该是你要推送的内容Future
    .throttle(1, 1.second, 1, ThrottleMode.Shaping)
    .map(TextMessage(_))

// 构建双向WebSocket Flow
val webSocketFlow: Flow[Message, Message, NotUsed] = Flow[Message]
  // 异步处理客户端消息,生成回复
  .mapAsync(parallelism = 1)(processClientMessage)
  // 合并主动推送的消息和回复消息,一起发给客户端
  .merge(activePushSource)

额外说明

  • 如果不需要主动推送,只需要回复客户端消息,那直接用Flow[Message].mapAsync(1)(processClientMessage)就够了
  • parallelism = 1可以根据你的并发需求调整,比如如果处理逻辑是线程安全的,可以设高一点
  • 要是需要更复杂的逻辑(比如根据客户端消息暂停主动推送,或者关联请求和回复),可以用Broadcast和Merge来拆分合并流,或者用Flow.fromSinkAndSourceCoupled来更精细控制双向流

这样你的WebSocket就能同时接收客户端消息并回复,还保留原来的主动推送功能啦!

内容的提问来源于stack exchange,提问作者S.K

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.28 06:23:13