自定义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的状态截图:
恳请提供调试思路或解决方案。
内容的提问来源于stack exchange,提问作者Bill'o
相关产品推荐
相关产品推荐

