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

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"))))
}

问题分析与解决方案

核心问题原因

  1. Source.actorRef生命周期冲突:Play返回响应时会负责materialize Source,但代码中在mapMaterializedValue里立即创建Actor并发送消息,Akka Stream 2.4.x中Source.actorRef的初始化未完成时接收消息,会导致内部VirtualPublisher注册失败,触发异常。
  2. 阻塞式IO与消息发送逻辑错误:StreamActor中用Future+Thread.sleep做IO读取,属于阻塞操作,违背Akka Stream异步非阻塞设计,容易引发消息积压和流状态异常。
  3. 流结束信号不正确: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 14:12:33