在前面的课程中,我们已经通过 AgentTaskManager 实现了单实例下的任务注册、并发控制、流式终止和资源释放。但那套实现只依赖本地 ConcurrentHashMap,一旦系统部署多个实例,问题就出现了,比如我现在启动两个实例,一个是8888端口,一个是8889端口,然后通过nginx进行转发,就可能出现,任务跑在 8888 端口上,停止请求却被 Nginx 转发到了 8889 端口,本地 map 里根本找不到这个任务,停止指令无处执行。本节课要解决的就是这个问题——如何让任务管理跨实例生效。改造后的 AgentTaskManager 引入了 Redis 分布式锁和 Pub/Sub 机制,使得任意实例都能正确停止运行在其他实例上的任务,同时保持了单实例场景下的简洁性和高性能。
从单实例到多实例:问题出在哪?
在真实的生产部署中,单个应用实例很难支撑所有用户请求。通常会通过 Nginx 等反向代理,将请求轮询分发到多个后端实例。我们的项目中也是这样做的:
## nginx
upstream agent_backend
server
server
}
server
listen
location
proxy_pass http
proxy_http_version
proxy_buffering off
proxy_cache off
proxy_set_header
proxy_read_timeout
}
这里有两个关键配置需要特别注意: proxy_read_timeout 是必须设置的 SSE 属于长连接,Nginx 默认 60 秒读超时,如果在该时间内没有数据传输,会强制断开连接,导致流式输出中断。 proxy_buffering off 建议开启 Nginx 默认会对响应进行缓冲,这可能导致后端生成的 token 被积攒后再批量发送,从而影响流式体验。虽然在某些情况下(如 token 输出频繁、数据量较小)即使不关闭缓冲也不会出现明显问题,但在高并发或复杂网络环境下仍可能导致延迟,因此建议显式关闭。 proxy_cache off 一般不是必须 SSE 响应通常带有 no-cache 头,Nginx 默认不会缓存,因此该配置更多是保险措施。 配置完成后,需要修改前端的接口配置:dodo-agent\src\main\resources\static\js\config.js
const
backendUrl
}
// 导出配置(用于非模块化环境)
window
在这种架构下,用户对话的请求可能被分发到不同的实例:
原有的 AgentTaskManager 只用 ConcurrentHashMap 做本地存储,8889 实例的 map 里根本没有这个任务,stopTask() 会直接返回 false,任务无法停止。
改造思路:Redis 分布式协作
要让多个实例协作管理任务,需要解决两个核心问题: 第一,任务注册的互斥性。同一个 conversationId 的任务只能在一个实例上运行,不能两个实例同时处理同一个会话。 第二,停止指令的路由。停止请求可能到达任意实例,但必须能精准传达给持有任务的那个实例。 改造方案如下:
┌──────────────────────────────────────────────────┐
│ Redis │
│ │
│ agent:task:{conversationId} → instanceId (SETNX) │
│ │
│ agent:stop (Pub/Sub Topic) │
│ │
└──────────┬───────────────────────┬───────────────┘
│ │
▼ ▼
┌──────────────┐ ┌──────────────┐
│ 实例 8888 │ │ 实例 8889 │
│ │ │ │
│ taskMap(本地) │ │ taskMap(本地) │
│ instanceId │ │ instanceId │
│ │ │ │
│ 订阅 stop │ │ 订阅 stop │
└──────────────┘ └──────────────┘
- Redis SETNX:任务注册时写入 agent:task:{conversationId} → instanceId,也就是agent:task:{conversationId} 为 key,instanceId为 value,保证同一会话只注册一次
- Redis Pub/Sub:停止请求通过 agent:stop 主题广播,所有实例监听,只有持有任务的实例执行停止,其他实例自动忽略即可。
任务注册
改造后的 registerTask 在本地检查之后,增加了 Redis 分布式锁:
public
log
log
taskMap
log
bucket.trySet() 是 Redis 的 SETNX 命令,同时设置了 30 分钟的 TTL。这样即使实例意外宕机,Redis key 也会自动过期,不会造成死锁。 注册流程如下:
实例 A 收到请求
│
├── 检查本地 taskMap → 无冲突
│
├── Redis SETNX → 成功(写入 instanceId)
│
└── 注册到本地 taskMap
---
实例 B 收到同一会话的请求
│
├── 检查本地 taskMap → 无冲突
│
├── Redis SETNX → 失败(key 已存在)
│
└── 拒绝注册,返回 null
hasRunningTask 也做了对应改造,除了检查本地 map,还会检查 Redis:
public
这样即使任务注册在另一个实例上,当前实例也能正确判断出该会话正在执行。
实例标识
每个 AgentTaskManager 实例在构造时生成一个 8 位随机 ID,作为整个分布式协调的身份标识:
public
log
instanceId 用于三个地方: - 注册时:写入 Redis value,标识任务的持有者 - 停止时:判断 Redis key 的持有者是否是本实例 - 续期时:校验 key 归属是否发生变化
停止会话:Pub/Sub
这部分是改造的核心,停止请求可能到达任意实例,但任务可能运行在其他实例上。我们需要一个机制把停止指令路由到正确的实例上。
stopTask 的四层判断
public
log
log
log
整个判断逻辑的流向:
stopTask(conversationId)
│
├── 本地 taskMap 有 → 直接停止
│
├── Redis key 不存在 → 没任务,返回 false
│
├── Redis 持有者是本实例 → 已在处理中,跳过
│
└── Redis 持有者是其他实例 → Pub/Sub 广播
Pub/Sub 订阅与远程停止
每个实例在启动时(afterPropertiesSet)都会订阅 agent:stop 主题:
public
listenerId
ttlRefreshScheduler
TTL_REFRESH_INTERVAL_MINUTES
TTL_REFRESH_INTERVAL_MINUTES
log
当某个实例发布广播后,所有实例都会收到消息,但只有持有任务的实例会处理:
private
log
完整的跨实例停止流程:
TTL 自动续期
任务注册时 Redis key 的 TTL 设为 30 分钟。但某些复杂任务(如深度研究)执行时间可能超过 10多分钟,如果 key 过期被删除,其他实例就可能重复注册同一会话的任务。 为此,每个实例启动了一个定时任务,每 5 分钟刷新本地所有任务的 TTL:
private
log
bucket
log
conversationId
taskMap
log
}
续期时会校验 holder 是否仍然是自己的 instanceId,防止误续别人的 key。
前端停止机制:abort 与 stopAgent
前端的 stopMessage 函数中同时使用了两种停止手段:
const
abortController
await APP_API
}
这两种方式的工作原理完全不同,理解它们的差异对于理解整个停止机制至关重要。
abort:基于 SSE 连接的直连停止
前端通过 fetch 建立 SSE 连接时,会创建一个 AbortController 并将其 signal 绑定到请求上:
abortController = new AbortController();
const response = await fetch(sseUrl, { signal: abortController.signal });
当调用 abortController.abort() 时,浏览器会立即断开这条 TCP 连接。关键在于这条连接是直接连到持有任务的那个后端实例的(中间经过 Nginx 转发)。所以后端检测到连接断开后,Reactive 流的 doOnCancel 回调会直接触发,进而调用 stopTask() 完成本地停止。
可以看到,当 abort() 生效时,整个 Pub/Sub 机制根本没有参与,8888 已经通过 doOnCancel 自己把自己停了,等 8889 收到 stopAgent 请求时,Redis key 都已经删了,直接返回 false。这意味着:只要 abort 生效,stopAgent 和 Pub/Sub 就完全是多余的。
abort 的局限
但光靠 abort 其实并不可靠,它的前提是必须要有SSE链接才可以。abort 依赖 SSE 连接的存在,比如:网络抖动,导致链接意外断开,abort() 同样无法触发。再比如如果系统中有其他组件需要主动停止某个任务(例如管理后台远程终止用户的执行任务),它并不持有用户的 SSE 连接,无法使用 abort,所以只能通过调用 stopAgent 接口配合 Pub/Sub 来完成跨实例停止。
验证 Pub/Sub:去掉 abort
由于 abort 会抢先完成停止,直接观察不到 Pub/Sub 的效果。如果想验证 Pub/Sub 机制是否正常工作,可以在前端代码中把 abort 注释掉:
const stopMessage = async () => {
if (!isSending.value) return;
// 注释掉 abort,让停止完全走后端 Pub/Sub 路径
// if (abortController) {
// abortController.abort();
// }
await APP_API.stopStream(backendUrl, currentChatId);
};
此时停止流程就会完全按照Pub/Sub 订阅与远程停止的时序图执行:stopAgent → Pub/Sub 广播 → handleRemoteStop → 停止任务。
建议:两者都保留
在实际项目中,建议两种机制都保留,各司其职:
用户点击停止
│
├── abort() → 快速路径:直接断 TCP,触发 doOnCancel,毫秒级响应
│
└── stopAgent API → 兜底路径:Pub/Sub 跨实例停止,覆盖 abort 失效的场景
abort 负责绝大多数正常场景下的快速停止,stopAgent + Pub/Sub 则为异常场景(标签页关闭、网络断开、管理后台远程停止等)提供可靠的兜底能力。
实例销毁时清理
当实例关闭(如重启、缩容)时,AgentTaskManager 通过 Spring 的 DisposableBean 回调执行清理:
@Override
public
stopTopic
log
ttlRefreshScheduler
log
}
第 3 步确保实例下线时,其持有的所有 Redis key 都被主动删除,而不是等待 30 分钟 TTL 自然过期。这样其他实例可以立即接管这些会话。
总结
从单实例到多实例的改造,核心变化是将 AgentTaskManager 的存储和协调机制从本地 ConcurrentHashMap 扩展到中间件 Redis。改造涉及三个关键机制: - Redis SETNX 分布式锁:保证同一会话的任务只能在一个实例上注册,解决了并发注册的互斥问题 - Redis Pub/Sub 广播:停止请求通过广播到达所有实例,由持有任务的实例执行停止,解决了跨实例路由问题 - TTL 自动续期:防止长任务在执行过程中 Redis key 过期导致的状态丢失 整个改造保持了单实例场景下的简洁性,本地任务直接走本地 map 快速路径,只有在涉及跨实例时才访问 Redis。