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

自定义Playframework BodyParser并发请求时Sink.asPublisher异常求助

自定义Playframework BodyParser并发请求异常排查求助

查过StackOverflow和GitHub相关问题,该问题在不同场景多次出现但无有效解决方案,特此求助。

我们使用标准自定义Playframework BodyParser将请求体转换为BinaryStream:

private val streamParser: BodyParser[AkkaStreams.BinaryStream] =
  BodyParser { _ =>
    Accumulator.source[ByteString].map(Right.apply)
  }

但在处理并发请求时会随机抛出java.lang.IllegalStateException,原因是有订阅者尝试订阅已被使用的Publisher。底层实现逻辑为创建Akka Stream的Sink,再生成仅支持单个订阅者的Publisher和Source:

def source[E]: Accumulator[E, Source[E, _]] = {
// If Akka streams ever provides Sink.source(), we should use that instead.
// https://github.com/akka/akka/issues/18406
    new SinkAccumulator(
      Sink
        .asPublisher[E](fanout = false)
        .mapMaterializedValue(publisher => 
      Future.successful(Source.fromPublisher(publisher)))
    )
  }

将fanout = true后会出现另一无有效信息的异常:

akka.stream.AbruptTerminationException: Processor actor [Actor[akka://application/system/Materializers/StreamSupervisor-0/flow-6131-1-fanoutPublisherSink#-1526957059]] terminated abruptly

修改Akka Streams配置仅能将异常转为警告或降低出现频率,无法触及根本原因:

akka {
  stream {
    materializer {
    debug-logging = on
    stage-errors-default-log-level = debug
    subscription-timeout {
      mode = warn
      timeout = 50s
    }
  }
}

关联的GitHub Issue已停滞多年,经多日排查,仍不清楚问题触发场景、是否存在共享状态或竞态条件,也无法稳定复现该问题。附上异常发生时Publisher的状态截图:
异常发生时Publisher状态

恳请提供调试思路或解决方案。

内容的提问来源于stack exchange,提问作者Bill'o

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 13:12:11