agentscope-ai / agentscope-ai/agentscope-java
[Bug]:Flux.create sink.complete() 延迟 30+ 秒才传播到下游
- Vorherrschende Sprache
- Java
- Sterne
- 5.6k
- Forks
- 1.3k
- Ø Merge
- 4 T. 12 Std.
- Gemergte PRs (30 T.)
- 77
Beschreibung
发现日期: 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` 之间的时间差
**注意**:偶发,可能需要多轮对话才触发。
---
Beitragsleitfaden
Bewertung
Dieses Issue wurde noch nicht bewertet.