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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 03:57:48