如何关闭Akka WebSocket服务器的消息缓冲功能?
问题分析与解决方案
嘿,我一眼就瞅出问题的根源了——你在Akka Streams的map操作里用了Thread.sleep(400),这可是阻塞流处理线程的大忌!
为啥会出现一分钟的缓冲延迟?
Akka Streams的调度器靠的是非阻塞异步操作来高效处理流元素。当你调用Thread.sleep时,直接把负责处理这个流的线程给堵死了——这个线程本来该把处理好的消息立刻推给客户端,结果被强制休眠,导致消息只能堆在缓冲区里,直到线程被释放后才会一股脑儿发出去。这就是为啥客户端要等一分钟才能收到批量消息的原因。
正确的实现姿势
要实现每隔400ms推送一行的需求,得用Akka Streams提供的非阻塞延迟操作,比如用Source.tick生成定时信号,再和文件流做zip。下面是修改后的完整代码:
object WebsocketServer extends App { implicit val actorSystem = ActorSystem("WebsocketServer") implicit val materializer = ActorMaterializer() implicit val executionContext = actorSystem.dispatcher val file = Paths.get("websocket-server/src/main/resources/EURUSD.txt") val fileSource = FileIO.fromPath(file) .via(Framing.delimiter(ByteString("\n"), Int.MaxValue)) .map(line => TextMessage(line.utf8String)) // 用定时信号和文件流绑定,实现间隔推送 val delayedSource: Source[TextMessage.Strict, Future[IOResult]] = Source.tick(0.millis, 400.millis, ()) .zip(fileSource) .map(_._2) // 只保留文件里的内容 .alsoTo(Sink.foreach(msg => println(msg.text))) // 异步打印推送内容,不阻塞流 def route = path("") { extractUpgradeToWebSocket { upgrade => complete(upgrade.handleMessagesWithSinkSource( Sink.ignore, delayedSource )) } } val bindingFuture = Http().bindAndHandle(route, "localhost", 8080) bindingFuture.onComplete { case Success(binding) => println(s"Server is listening on ws://localhost:8080") case Failure(e) => println(s"Binding failed with ${e.getMessage}") actorSystem.terminate() } }
关键改动说明
- 删掉了
map里的Thread.sleep和打印:换成alsoTo(Sink.foreach)来异步打印消息,彻底避免阻塞流处理线程。 - 用
Source.tick生成定时信号:这个源会从启动开始,每隔400ms发送一个空信号,和文件流zip之后,文件里的每一行都会等这个信号才会往下走,完美实现每隔400ms推送一行的效果,全程非阻塞。 - 如果你用的是Akka 2.6及以上版本,
ActorMaterializer已经被弃用,可以换成SystemMaterializer:import akka.stream.SystemMaterializer implicit val materializer = SystemMaterializer(actorSystem).materializer
这样修改后,客户端就能每隔400ms收到一条消息,不会再出现批量延迟的情况啦!
内容的提问来源于stack exchange,提问作者Tom Lous
相关产品推荐
相关产品推荐

