如何从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
相关产品推荐
相关产品推荐

