Play Framework返回Akka Source分块音频时触发IllegalStateException问题
Play Framework 2.5.18结合Akka Stream 2.4.20实现MP3分块传输时抛出IllegalStateException
问题描述
使用Play Framework 2.5.18结合Akka Stream 2.4.20实现MP3音频分块传输,创建Akka Source并通过Ok.chunked返回结果时抛出java.lang.IllegalStateException。改用HttpEntity.Streamed方式返回,仍出现相同异常。
错误日志
[ERROR 2023-08-02 00:13:13,223] [application-akka.actor.default-dispatcher-32 c.x.a.c.DefaultErrorHandler] - A server error occurred for speech-staging.ai.srv:speech-staging.ai.srv/stream_audio_v2,from: 10.38.162.91 play.core.server.common.ServerResultException: Error converting Play Result for server backend at play.core.server.common.ServerResultUtils$.play$core$server$common$ServerResultUtils$$handleConversionError$1(ServerResultUtils.scala:108) at play.core.server.common.ServerResultUtils$.resultConversionWithErrorHandling(ServerResultUtils.scala:129) at play.core.server.netty.NettyModelConversion.convertResult(NettyModelConversion.scala:251) at play.core.server.netty.PlayRequestHandler$$anonfun$play$core$server$netty$PlayRequestHandler$$handleAction$2$$anonfun$apply$4$$anonfun$apply$5.apply(PlayRequestHandler.scala:273) at play.core.server.netty.PlayRequestHandler$$anonfun$play$core$server$netty$PlayRequestHandler$$handleAction$2$$anonfun$apply$4$$anonfun$apply$5.apply(PlayRequestHandler.scala:267) at scala.concurrent.Future$$anonfun$flatMap$1.apply(Future.scala:253) at scala.concurrent.Future$$anonfun$flatMap$1.apply(Future.scala:251) at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:36) at play.api.libs.iteratee.Execution$trampoline$.executeScheduled(Execution.scala:109) at play.api.libs.iteratee.Execution$trampoline$.execute(Execution.scala:71) at scala.concurrent.impl.CallbackRunnable.executeWithValue(Promise.scala:44) at scala.concurrent.impl.Promise$DefaultPromise.tryComplete(Promise.scala:252) at scala.concurrent.Promise$class.complete(Promise.scala:55) at scala.concurrent.impl.Promise$DefaultPromise.complete(Promise.scala:157) at scala.concurrent.Future$$anonfun$flatMap$1$$anonfun$apply$3.apply(Future.scala:256) at scala.concurrent.Future$$anonfun$flatMap$1$$anonfun$apply$3.apply(Future.scala:256) at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:36) at scala.concurrent.BatchingExecutor$Batch$$anonfun$run$1.processBatch$1(BatchingExecutor.scala:63) at scala.concurrent.BatchingExecutor$Batch$$anonfun$run$1.apply$mcV$sp(BatchingExecutor.scala:78) at scala.concurrent.BatchingExecutor$Batch$$anonfun$run$1.apply(BatchingExecutor.scala:55) at scala.concurrent.BatchingExecutor$Batch$$anonfun$run$1.apply(BatchingExecutor.scala:55) at scala.concurrent.BlockContext$.withBlockContext(BlockContext.scala:72) at scala.concurrent.BatchingExecutor$Batch.run(BatchingExecutor.scala:54) at scala.concurrent.Future$InternalCallbackExecutor$.unbatchedExecute(Future.scala:601) at scala.concurrent.BatchingExecutor$class.execute(BatchingExecutor.scala:106) at scala.concurrent.Future$InternalCallbackExecutor$.execute(Future.scala:599) at scala.concurrent.impl.CallbackRunnable.executeWithValue(Promise.scala:44) at scala.concurrent.impl.Promise$KeptPromise.onComplete(Promise.scala:337) at scala.concurrent.Future$$anonfun$flatMap$1.apply(Future.scala:256) at scala.concurrent.Future$$anonfun$flatMap$1.apply(Future.scala:251) at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:36) at akka.dispatch.BatchingExecutor$AbstractBatch.processBatch(BatchingExecutor.scala:55) at akka.dispatch.BatchingExecutor$BlockableBatch$$anonfun$run$1.apply$mcV$sp(BatchingExecutor.scala:91) at akka.dispatch.BatchingExecutor$BlockableBatch$$anonfun$run$1.apply(BatchingExecutor.scala:91) at akka.dispatch.BatchingExecutor$BlockableBatch$$anonfun$run$1.apply(BatchingExecutor.scala:91) at scala.concurrent.Future$$anonfun$flatMap$1.apply(Future.scala:256) at scala.concurrent.Future$$anonfun$flatMap$1.apply(Future.scala:251) at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:36) at akka.dispatch.BatchingExecutor$AbstractBatch.processBatch(BatchingExecutor.scala:55) at akka.dispatch.BatchingExecutor$BlockableBatch$$anonfun$run$1.apply$mcV$sp(BatchingExecutor.scala:91) at akka.dispatch.BatchingExecutor$BlockableBatch$$anonfun$run$1.apply(BatchingExecutor.scala:91) at akka.dispatch.BatchingExecutor$BlockableBatch$$anonfun$run$1.apply(BatchingExecutor.scala:91) at scala.concurrent.BlockContext$.withBlockContext(BlockContext.scala:72) at akka.dispatch.BatchingExecutor$BlockableBatch.run(BatchingExecutor.scala:90) at akka.dispatch.TaskInvocation.run(AbstractDispatcher.scala:39) at akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(AbstractDispatcher.scala:415) at scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) at scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339) at scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) at scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107) Caused by: java.lang.IllegalStateException: internal error at akka.stream.impl.VirtualPublisher.registerPublisher(StreamLayout.scala:832) at akka.stream.impl.MaterializerSession.akka$stream$impl$MaterializerSession$$doSubscribe(StreamLayout.scala:1034) at akka.stream.impl.MaterializerSession$$anonfun$materialize$3$$anonfun$apply$3.apply(StreamLayout.scala:925) at akka.stream.impl.MaterializerSession$$anonfun$materialize$3$$anonfun$apply$3.apply(StreamLayout.scala:924) at scala.collection.Iterator$class.foreach(Iterator.scala:891) at scala.collection.AbstractIterator.foreach(Iterator.scala:1334) at akka.stream.impl.MaterializerSession$$anonfun$materialize$3.apply(StreamLayout.scala:924) at akka.stream.impl.MaterializerSession$$anonfun$materialize$3.apply(StreamLayout.scala:924) at scala.collection.immutable.List.foreach(List.scala:392) at akka.stream.impl.MaterializerSession.materialize(StreamLayout.scala:924) at akka.stream.impl.ActorMaterializerImpl.materialize(ActorMaterializerImpl.scala:256) at akka.stream.impl.ActorMaterializerImpl.materialize(ActorMaterializerImpl.scala:146) at akka.stream.scaladsl.RunnableGraph.run(Flow.scala:350) at akka.stream.scaladsl.Source.runWith(Source.scala:81) at play.core.server.netty.NettyModelConversion.play$core$server$netty$NettyModelConversion$$createChunkedResponse(NettyModelConversion.scala:272) at play.core.server.netty.NettyModelConversion$$anonfun$convertResult$1.apply(NettyModelConversion.scala:205) at play.core.server.netty.NettyModelConversion$$anonfun$convertResult$1.apply(NettyModelConversion.scala:182) at play.core.server.common.ServerResultUtils$.resultConversionWithErrorHandling(ServerResultUtils.scala:127) ... 41 common frames omitted
核心代码
StreamTest控制器代码
class StreamTest @Inject()(actorSystem: ActorSystem)(implicit exec: ExecutionContext) extends Controller { def stream_audio_v2(): Action[AnyContent] = Action.async { req => val source = Source.actorRef[ByteString](10000, OverflowStrategy.fail).mapMaterializedValue { sourceActor => val actor = actorSystem.actorOf(Props(new StreamActor(sourceActor)).withDispatcher("assist-dispatcher")) actor ! "stream_audio_v2" } Future.successful(Ok.chunked(source).as("audio/mp3")) } }
StreamActor代码
class StreamActor(outActor: ActorRef = null) extends Actor { override def receive: Receive = { case "stream_audio_v2" => log.info("receive stream_audio_v2") val out = outActor log.info("out put streamProducer.stream") Future { val source = StreamActor.getSource() var buffer = new Array[Byte](1024) var length = 0 while (length != -1) { length = source.read(buffer) out ! ByteString(buffer.clone()) Thread.sleep(1) } source.close() out ! akka.actor.Status.Success(()) } } }
尝试过的解决方案
改用HttpEntity.Streamed方式构造响应,代码如下,但异常依旧:
def stream_audio_v2_5(): Action[AnyContent] = Action.async { req => val source = Source.actorRef[ByteString](10000, OverflowStrategy.fail).mapMaterializedValue { sourceActor => val actor = actorSystem.actorOf(Props(new StreamActor(sourceActor)).withDispatcher("assist-dispatcher")) actor ! "stream_audio_v2" } Future.successful(Ok.sendEntity(HttpEntity.Streamed(source, None, Some("audio/mp3")))) }
问题分析与解决方案
核心问题原因
- Source.actorRef生命周期冲突:Play返回响应时会负责materialize Source,但代码中在
mapMaterializedValue里立即创建Actor并发送消息,Akka Stream 2.4.x中Source.actorRef的初始化未完成时接收消息,会导致内部VirtualPublisher注册失败,触发异常。 - 阻塞式IO与消息发送逻辑错误:
StreamActor中用Future+Thread.sleep做IO读取,属于阻塞操作,违背Akka Stream异步非阻塞设计,容易引发消息积压和流状态异常。 - 流结束信号不正确:
Source.actorRef需要Complete或Failure信号终止流,但代码发送的是akka.actor.Status.Success(()),导致流无法正常收尾,引发内部状态错误。
修复方案
方案1:用Akka Stream原生IO源替代Actor手动读取
直接使用Akka Stream的IO工具类处理音频读取,避免Actor中间层:
class StreamTest @Inject()(actorSystem: ActorSystem)(implicit mat: Materializer) extends Controller { def stream_audio_v2(): Action[AnyContent] = Action.async { req => // 替换为Akka Stream的InputStream Source,模拟原有的延迟发送 val audioSource = Source.fromInputStream(() => StreamActor.getSource()) .map(ByteString(_)) .throttle(1, 1.milliseconds, 1024, ThrottleMode.shaping) Future.successful(Ok.chunked(audioSource).as("audio/mp3")) } }
方案2:修复Source.actorRef的使用逻辑
若必须保留Actor,调整消息发送时机并使用正确的流结束信号:
class StreamTest @Inject()(actorSystem: ActorSystem)(implicit mat: Materializer) extends Controller { def stream_audio_v2(): Action[AnyContent] = Action.async { req => val source = Source.actorRef[ByteString](10000, OverflowStrategy.fail) .watchTermination() { (_, termination) => termination.onComplete(_ => actorSystem.stop(actor)) actor } .mapMaterializedValue { sourceActor => val actor = actorSystem.actorOf(Props(new StreamActor(sourceActor)).withDispatcher("assist-dispatcher")) // 延迟发送消息,确保Source初始化完成 actorSystem.scheduler.scheduleOnce(10.milliseconds, actor, "stream_audio_v2") actor } Future.successful(Ok.chunked(source).as("audio/mp3")) } } // 修改StreamActor的读取逻辑与结束信号 class StreamActor(outActor: ActorRef = null) extends Actor { private var source: InputStream = _ private val buffer = new Array[Byte](1024) override def receive: Receive = { case "stream_audio_v2" => log.info("receive stream_audio_v2") source = StreamActor.getSource() // 用Akka调度器替代Thread.sleep,避免阻塞 context.system.scheduler.schedule(0.milliseconds, 1.milliseconds, self, "read_chunk") case "read_chunk" => val length = source.read(buffer) if (length != -1) { outActor ! ByteString(buffer.take(length)) // 只发送有效字节 } else { source.close() outActor ! akka.stream.scaladsl.Source.Complete // 发送正确的流结束信号 context.stop(self) } } }
方案3:升级Akka Stream版本
Akka Stream 2.4.x存在流materialize相关的已知bug,升级到2.4.21及以上兼容Play 2.5的版本,可修复部分内部状态异常问题。
内容的提问来源于stack exchange,提问作者CannyMiao
相关产品推荐
相关产品推荐

