Scala Akka使用Actor作为WebSocket客户端流数据源与Sink的编译问题
问题定位
- 输入输出类型不匹配:WebSocket客户端Flow要求输入为
Message类型(即TextMessage/BinaryMessage子类),你之前声明的ActorSource产出类型为Any或String,与要求的输入类型不符,触发最早的类型不匹配错误。 - 偏函数类型推断失败:
completionMatcher是偏函数,没有显式指定输入类型时编译器无法自动推断,导致整个流的物化值类型全部丢失,被识别为Any,后续调用flatMap、!等方法全部报错。 - 重复物化Source:第一次尝试的代码中对同一个
actorSource调用了两次run(),Akka Stream的每个Source只能被物化一次,第二次运行的流会直接失效。 - 发送消息类型冲突:更新后的代码声明ActorSource产出
String,但实际给Actor发送的是TextMessage实例,类型不匹配。
正确实现方案
完整可运行代码如下,已经包含将接收消息转发到其他Actor的逻辑:
import akka.{Done, NotUsed} import akka.actor.{Actor, ActorRef, ActorSystem, Props} import akka.http.scaladsl.Http import akka.http.scaladsl.model.ws.{Message, TextMessage, WebSocketRequest, WebSocketUpgradeResponse} import akka.stream.{CompletionStrategy, OverflowStrategy} import akka.stream.scaladsl.{Flow, Keep, Sink, Source} import akka.http.scaladsl.model.StatusCodes import scala.concurrent.Future object WsClient { def main(args: Array[String]): Unit = { implicit val system = ActorSystem("WsClientSystem") import system.dispatcher // 接收WebSocket消息的目标Actor,可替换为你自己的业务Actor val messageReceiverActor: ActorRef = system.actorOf(Props[MessageReceiverActor]) // WebSocket接收端Sink,所有消息直接转发到指定Actor val incoming: Sink[Message, Future[Done]] = Sink.actorRef( ref = messageReceiverActor, onCompleteMessage = Done, onFailureMessage = ex => new RuntimeException(s"流异常: ${ex.getMessage}") ) // 声明ActorSource,产出类型与WebSocket输入要求对齐 val actorSource: Source[TextMessage.Strict, ActorRef] = Source.actorRef[TextMessage.Strict]( // 显式指定偏函数类型,解决类型推断失败问题 completionMatcher = PartialFunction[Any, CompletionStrategy] { case Done => CompletionStrategy.immediately }, failureMatcher = PartialFunction.empty, bufferSize = 100, overflowStrategy = OverflowStrategy.dropHead ) val webSocketFlow: Flow[Message, Message, Future[WebSocketUpgradeResponse]] = Http().webSocketClientFlow( WebSocketRequest("wss://socket.polygon.io/stocks") ) // 组合流并保留所有需要的物化值 val ((sendActor: ActorRef, upgradeResponse: Future[WebSocketUpgradeResponse]), closed: Future[Done]) = actorSource .viaMat(webSocketFlow)(Keep.both) .toMat(incoming)(Keep.both) .run() // 连接状态校验 val connected: Future[Done] = upgradeResponse.flatMap { upgrade => if (upgrade.response.status == StatusCodes.SwitchingProtocols) { Future.successful(Done) } else { throw new RuntimeException(s"连接失败: ${upgrade.response.status}") } } // 连接成功后再发送认证、订阅消息,避免消息丢失 connected.onComplete { connResult => println(s"连接结果: $connResult") if (connResult.isSuccess) { sendActor ! TextMessage.Strict("""{"auth":"APIKEY-123"}""") sendActor ! TextMessage.Strict("""{"topic":"topic123"}""") } } closed.onComplete { result => println(s"连接关闭: $result") } } } // 示例:处理WebSocket接收消息的Actor实现 class MessageReceiverActor extends Actor { override def receive: Receive = { case TextMessage.Strict(text) => println(s"收到消息: $text") // 此处写你自己的业务处理逻辑 case Done => println("WebSocket流已正常结束") case ex: RuntimeException => println(s"WebSocket流异常: ${ex.getMessage}") } }
核心改动说明
- 给
completionMatcher显式指定偏函数类型PartialFunction[Any, CompletionStrategy],解决编译器类型推断失败的问题,流的物化值类型可以正确传递,不再被识别为Any。 - 将ActorSource的产出类型明确指定为
TextMessage.Strict,和WebSocket Flow的输入类型完全对齐,不需要额外类型转换,解决最早的类型不匹配错误。 - 替换原来的打印Sink为
Sink.actorRef,所有收到的WebSocket消息都会直接转发到指定Actor,满足消息转发需求。 - 将发送消息的逻辑放到连接成功的回调中,避免连接未建立就发消息导致消息丢失。
内容的提问来源于stack exchange,提问作者Blankman
相关产品推荐
相关产品推荐

