Akka Streams向QueueSource发送元素时出现StreamDetachedException问题
跨域文件下载流中QueueSource触发StreamDetachedException问题排查
我们基于Akka Streams和Akka Http实现了DownloadFileFlow类用于跨域文件下载,但约50%的概率下,向QueueSource发送元素时会立即触发StreamDetachedException错误,堆栈信息如下:
java.util.concurrent.CompletionException: akka.stream.StreamDetachedException: Stage with GraphStageLogic akka.stream.impl.QueueSource$$anon$1-queueSource stopped before async invocation was processed at java.base/java.util.concurrent.CompletableFuture.encodeThrowable(CompletableFuture.java:332) at java.base/java.util.concurrent.CompletableFuture.uniAcceptNow(CompletableFuture.java:747) at java.base/java.util.concurrent.CompletableFuture.uniAcceptStage(CompletableFuture.java:735) at java.base/java.util.concurrent.CompletableFuture.thenAcceptAsync(CompletableFuture.java:2186) at scala.concurrent.java8.FuturesConvertersImpl$CF.thenAccept(FutureConvertersImpl.scala:29) at scala.concurrent.java8.FuturesConvertersImpl$CF.thenAccept(FutureConvertersImpl.scala:18) at x.y.z.DownloadFileFlow.offer(DownloadFileFlow.java:60)
我们注意到akka.stream.impl.QueueSource中存在postStop方法,该方法会抛出上述异常,但目前不清楚流停止的触发原因,怀疑是否由资源泄漏导致(例如某些场景下未消费Http响应实体)。
相关代码及配置
1. SourceQueue引用
AtomicReference<SourceQueueWithComplete<FileDownloadEnvelope>> queueRef = new AtomicReference<>()
2. 流定义
RestartSource.onFailuresWithBackoff( RestartSettings.create(props.getRecoverMinBackOff(), props.getRecoverMaxBackOff(), RANDOM_FACTOR), () -> Source.<FileDownloadEnvelope>queue(props.getBufferSize(), OverflowStrategy.backpressure(), props.getMaxConcurrentOffers()) .filter(this::checkExpiration) .map(x -> Pair.create(createContext(x), x)) .flatMapConcat(this::authIfNeeded) .mapAsyncUnordered(props.getParallelism(), this::process) .map(this::reply) .mapMaterializedValue(x -> { this.queueRef.set(x); return x; })) .runWith(Sink.ignore(), actorSystem);
3. 发送元素的offer方法
public void offer(FileDownloadEnvelope envelope) { queueRef.get().offer(envelope) .thenAccept(x -> { if (!x.isEnqueued()) replyError(envelope); }).exceptionally(e -> { // StreamDetachedException catched here replyError(envelope); return null; }); }
Akka Http配置
akka.http.host-connection-pool.min-connections: "10" akka.http.host-connection-pool.max-connections: "2000" akka.http.host-connection-pool.max-open-requests: "4096" akka.http.host-connection-pool.max-retries: "0" akka.http.host-connection-pool.client.connecting-timeout: 2s akka.http.host-connection-pool.client.idle-timeout: 5s
流配置
BUFFER-SIZE: "2000" MAX-CONCURRENT-OFFERS: "2000" PARALLELISM: "70" RECOVER-MIN-BACK-OFF: 100ms RECOVER-MAX-BACK-OFF: 500ms
使用版本
- akka & akka-stream: 2.6.19
- akka-http: 10.2.0
内容的提问来源于stack exchange,提问作者Arya
相关产品推荐
相关产品推荐

