Skip to content

[Bug]:Flux.create sink.complete() 延迟 30+ 秒才传播到下游 #2279

Description

@jiajingshan

发现日期: 2026-07-18
偶尔发生,但 UX 差:前端卡在"生成中"状态 30s


1. 环境

AgentScope 版本 2.0.0
JDK 21
模型后端 SGLang (http://10.1.4.20:20010)
模型 qwen3.6-27b
思考模式 enable_thinking=true
传输层 JDK HttpClient HTTP/1.1

2. 现象

使用 HarnessAgent.streamEvents(msg, rc) 进行流式对话。模型输出在 00:15:33 完成(AgentTraceMiddleware POST_REASONING 日志出现),但 sink.complete() 的信号延迟到 00:16:03 才到达下游——中间 30 秒无任何事件,前端一直显示"生成中"。

日志时序

00:15:33 [boundedElastic-4] POST_REASONING text: ...   ← 模型输出完毕
00:15:33 [boundedElastic-4] 所有 Middleware onComplete
                         ═══════ 30 秒空白 ═══════
00:16:03 [boundedElastic-5] POST_CALL response          ← Agent 最终完成

触发规律

  • 偶尔发生(非必现)
  • 开启思考模式(enable_thinking=true)时更常见
  • 与模型、后端无关(SGLang/vLLM 均出现过)

3. 根因定位

3.1 问题代码路径

ReActAgent.buildAgentStream()(agentscope-core:2.0.0)使用 Flux.create 桥接内部 Mono 管道与外部事件流:

Flux.create(sink -> {
    AgentStartEvent start = new AgentStartEvent(...);
    sink.next(start);
    
    Mono<Msg> mono = runLifecycle(input, middlewareChain);
    mono.contextWrite(ctx -> ctx.put("eventSink", sink))
        .doFinally(sig -> {
            sink.next(new AgentEndEvent(agentName));
            sink.complete();  // ← 此处偶尔延迟 30s
        })
        .subscribe(
            msg -> sink.next(new AgentResultEvent(msg)),
            err -> sink.error(err)
        );
    sink.onCancel(subscription);
}, OverflowStrategy.BUFFER);  // ★ BUFFER 模式可能加剧延迟

3.2 延迟机制

  1. runLifecycle 的 Mono 在 boundedElastic-4 线程完成
  2. doFinally 回调在 boundedElastic-4 上触发
  3. sink.complete() 调用 FluxSink.complete(),但信号通过 BUFFER 策略传播到下游
  4. 下游引用了 concatMap 等操作符,这些操作符可能在不同的 boundedElastic 线程上处理
  5. 信号从完成线程传播到下游消费者线程的调度延迟达到 30 秒(当线程池繁忙或有其他背压条件时)

3.3 为什么 BUFFER

OverflowStrategy.BUFFER 意味着未消费的事件在内存中缓冲。当下游消费慢(如 SSE 序列化 + HTTP 写回),缓冲累积,sink.complete() 信号排在缓冲队列尾部。


4. 建议修复

方案 A(推荐):doFinally 中用 tryEmitComplete

.doFinally(sig -> {
    sink.next(new AgentEndEvent(agentName));
    EmitResult result = sink.tryEmitComplete();
    if (result.isFailure()) {
        log.warn("sink.tryEmitComplete failed: {}", result);
    }
})

tryEmitComplete() 是非阻塞 API,失败时不阻塞调用线程,由 Reactor 内部重试机制兜底。

方案 B:Flux.push + LATEST

Flux.push(sink -> { ... }, OverflowStrategy.LATEST)

Flux.push 是异步信号发射器,LATEST 策略丢弃未被消费的旧事件——相比 BUFFER 不会累积积压队列。但这可能丢失事件,需要权衡。

方案 C:分离 sink.complete() 到独立调度

.doFinally(sig -> {
    sink.next(new AgentEndEvent(agentName));
    // 在无背压的调度器上完成
    Schedulers.single().schedule(() -> sink.complete());
})

5. 复现方法

  1. HarnessAgent.streamEvents() 进行流式对话
  2. 开启思考模式(enable_thinking=true
  3. 使用 qwen3.6-27b 这样的中等规模模型(输出 token 较多)
  4. 连续发送多条长回复请求
  5. 观察 POST_REASONINGPOST_CALL 之间的时间差

注意:偶发,可能需要多轮对话才触发。


Metadata

Metadata

Assignees

No one assigned

    Labels

    area/core/agentAgent runtime, pipeline, hooks, planbugSomething isn't working

    Type

    No type

    Projects

    Status
    Backlog

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions