agentscope-ai / agentscope-ai/agentscope-java

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

Aperta
#2,279 3 commenti 0 reazioni 0 assegnatari Vedi su GitHub
area/core/agent bug
Lingua principale
Java
Stelle
5.6k
Fork
1.3k
Merge medio
4g 12h
PR unite (30g)
77

Descrizione

发现日期: 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 管道与外部事件流:

```java
Flux.create(sink -> {
AgentStartEvent start = new AgentStartEvent(...);
sink.next(start);

Mono 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`

```java
.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`

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

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

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

```java
.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_REASONING` 和 `POST_CALL` 之间的时间差

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

---

Guida per i contributori

Apri la guida per i contributori

Valutazione

Questa issue non è ancora stata valutata.

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.