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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 04:27:03