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

Akka Stream客户端遇超时及"This publisher only supports one subscriber"错误求助

Akka Stream客户端三类异常:超时、死锁/瓶颈问题排查与解决

在使用Akka Stream的客户端代码中,常会遇到调用失败、超时甚至死锁/瓶颈情况,但服务端运行完全正常。这类问题对应三类典型异常,以下是具体排查与解决方案:


异常分类及处理

1. java.lang.IllegalStateException: This publisher only supports one subscriber

触发原因:
底层的HandlerPublisher(来自Play WS依赖)是单订阅者模型的Publisher,同一个实例仅允许被订阅一次。常见场景:

  • 重试逻辑中复用了同一个WS响应的BodySource或Publisher实例,多次发起订阅
  • 同一个流实例被传给多个Sink处理
  • 流Graph中错误共享了非线程安全的Publisher组件

解决方法:

  • 每次处理响应流时,重新发起WS请求,不复用之前的响应流实例
  • 重试逻辑中,确保每次重试都创建全新的流Graph和请求实例
  • 禁止在多线程或流分支中共享同一个WS响应的Source对象

对应堆栈:

java.io.IOException: java.lang.IllegalStateException: This publisher only supports one subscriber
    at akka.stream.impl.io.InputStreamAdapter.$anonfun$read$5(InputStreamSinkStage.scala:170)
    at apply @ com.mycompany.utils.Retry.$anonfun$run$2(Retry.scala:29)
    at fromFuture @ com.mycompany.utils.Retry.$anonfun$run$2(Retry.scala:29)
    at flatMap @ retry.package$RetryingOnSomeErrorsPartiallyApplied.$anonfun$apply$3(package.scala:101)
    at tailRecM @ retry.package$RetryingOnSomeErrorsPartiallyApplied.apply(package.scala:100)
    at *> @ retry.package$.$anonfun$retryingOnSomeErrorsImpl$3(package.scala:79)
    at map @ retry.package$.$anonfun$retryingOnSomeErrorsImpl$3(package.scala:76)
    at apply @ com.mycompany.utils.Retry.logError(Retry.scala:51)
Caused by: java.lang.IllegalStateException: This publisher only supports one subscriber
    at play.shaded.ahc.com.typesafe.netty.HandlerPublisher.subscribe(HandlerPublisher.java:167)
    at play.api.libs.ws.ahc.StandaloneAhcWSClient$$anon$2.subscribe(StandaloneAhcWSClient.scala:120)
    at akka.stream.impl.fusing.ActorGraphInterpreter$BatchingActorInputBoundary.preStart(ActorGraphInterpreter.scala:148)
    at akka.stream.impl.fusing.GraphInterpreter.init(GraphInterpreter.scala:306)
    at akka.stream.impl.fusing.GraphInterpreterShell.init(ActorGraphInterpreter.scala:619)
    at akka.stream.impl.fusing.ActorGraphInterpreter.tryInit(ActorGraphInterpreter.scala:727)
    at akka.stream.impl.fusing.ActorGraphInterpreter.preStart(ActorGraphInterpreter.scala:776)
    at akka.actor.Actor.aroundPreStart(Actor.scala:548)
    at akka.actor.Actor.aroundPreStart$(Actor.scala:548)
    at akka.stream.impl.fusing.ActorGraphInterpreter.aroundPreStart(ActorGraphInterpreter.scala:716)
    at akka.actor.ActorCell.create(ActorCell.scala:644)
    at akka.actor.ActorCell.invokeAll$1(ActorCell.scala:514)
    at akka.actor.ActorCell.systemInvoke(ActorCell.scala:536)
    at akka.dispatch.Mailbox.processAllSystemMessages(Mailbox.scala:295)
    at akka.dispatch.Mailbox.run(Mailbox.scala:230)
    at akka.dispatch.Mailbox.exec(Mailbox.scala:243)
    at java.base/java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:290)
    at java.base/java.util.concurrent.ForkJoinPool$WorkQueue.topLevelExec(ForkJoinPool.java:1020)
    at java.base/java.util.concurrent.ForkJoinPool.scan(ForkJoinPool.java:1656)
    at java.base/java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1594)
    at java.base/java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:183)

2. java.io.IOException: Timeout after 60 seconds waiting for Initialized message from stage

触发原因:
Akka Stream的InputStreamAdapter组件在等待流阶段初始化信号时超时。常见场景:

  • 流Graph的上游Source未正确启动,无法发送初始化信号
  • 客户端线程池被阻塞操作占用,导致ActorGraphInterpreter无法处理初始化事件
  • WS请求连接建立后,服务端未及时响应,导致流初始化延迟

解决方法:

  • 检查流Graph构建逻辑,确保上游Source能正常发射元素,无意外阻塞
  • 调整InputStreamAdapter初始化超时时间,在配置文件中添加:
    akka.stream.input-stream-adapter.initialized-timeout = 120s
    
  • 将流处理中的阻塞代码移到专门的阻塞线程池(如akka.stream.blocking-io-dispatcher)

对应堆栈:

java.io.IOException: Timeout after 60 seconds waiting for Initialized message from stage
    at akka.stream.impl.io.InputStreamAdapter.waitIfNotInitialized(InputStreamSinkStage.scala:224)
    at akka.stream.impl.io.InputStreamAdapter.executeIfNotClosed(InputStreamSinkStage.scala:130)
    at akka.stream.impl.io.InputStreamAdapter.read(InputStreamSinkStage.scala:156)
    at com.fasterxml.jackson.core.json.ByteSourceJsonBootstrapper.ensureLoaded(ByteSourceJsonBootstrapper.java:539)
    at com.fasterxml.jackson.core.json.ByteSourceJsonBootstrapper.detectEncoding(ByteSourceJsonBootstrapper.java:133)
    at com.fasterxml.jackson.core.json.ByteSourceJsonBootstrapper.constructParser(ByteSourceJsonBootstrapper.java:256)
    at com.fasterxml.jackson.core.JsonFactory._createParser(JsonFactory.java:1655)
    at com.fasterxml.jackson.core.JsonFactory.createParser(JsonFactory.java:1083)
    at play.api.libs.json.jackson.JacksonJson$.parseJsValue(JacksonJson.scala:285)
    at play.api.libs.json.StaticBinding$.parseJsValue(StaticBinding.scala:21)
    at play.api.libs.json.Json$.parse(Json.scala:175)

3. java.io.IOException: Timeout on waiting for new data

触发原因:
InputStreamAdapter在等待流的新数据时超时,说明上游长时间未发射元素。常见场景:

  • 服务端响应缓慢,未及时返回数据
  • 流背压策略配置不合理,导致上游被过度压制无法发送数据
  • 重试逻辑中未正确关闭之前的流实例,导致资源泄漏,新流无法获取数据

解决方法:

  • 调整InputStreamAdapter数据读取超时时间,在配置文件中添加:
    akka.stream.input-stream-adapter.read-timeout = 120s
    
  • 检查服务端响应性能,优化接口或增加客户端超时阈值
  • 重试逻辑中,每次重试前关闭之前的流资源,避免资源占用
  • 调整流背压配置,使用buffer或conflate操作缓解背压问题

对应堆栈:

java.io.IOException: Timeout on waiting for new data
    at akka.stream.impl.io.InputStreamAdapter.$anonfun$read$5(InputStreamSinkStage.scala:171)
    at apply @ com.mycompany.utils.Retry.$anonfun$run$2(Retry.scala:29)
    at fromFuture @ com.mycompany.utils.Retry.$anonfun$run$2(Retry.scala:29)
    at flatMap @ retry.package$RetryingOnSomeErrorsPartiallyApplied.$anonfun$apply$3(package.scala:101)
    at tailRecM @ retry.package$RetryingOnSomeErrorsPartiallyApplied.apply(package.scala:100)
    at *> @ retry.package$.$anonfun$retryingOnSomeErrorsImpl$3(package.scala:79)
    at map @ retry.package$.$anonfun$retryingOnSomeErrorsImpl$3(package.scala:76)
    at apply @ com.mycompany.utils.Retry.logError(Retry.scala:51)

内容的提问来源于stack exchange,提问作者Gaël J

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 03:22:02