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
相关产品推荐
相关产品推荐

