concurrency-issues-and-fixes.md 45 KB

多用户并发问题分析与修改方案(ai-server)

审计日期:2026-09-30 | 修复实施:2026-09-30(见文末第九节) 范围:ai-server/src/main/java 全量(约 800 个 Java 文件) 方法:3 轮全仓静态扫描(共享可变状态 / 业务竞态与数据隔离 / 线程池与锁)+ 盲区补查(Redis / 全文检索 / 事务边界),全部 P0 与 P1 关键行号经人工逐条核实。 结论计数:P0 高危 7 项 | P1 中危 12 项 | P2 低危 10 项,另有附录列出已排查为安全的部分(避免重复排查)。


0. 结论摘要

项目已有一套较完整的并发基础设施(CaseContextHolder ScopedValue 上下文传播、CaseDataCache 按 caseId 分桶、研判流双保险锁、附件 CAS 防重入等,详见附录 A/B),最近也修过附件上传竞态。但"数据清洗完成回调"这一条链路集中了 3 个高危问题(全局静态人员缓存跨案件串写、删除全局临时目录、全局模板缓存无同步重建),多用户同时清洗/治理时必现串数据、丢数据或偶发崩溃,是本次审计的最高优先级。

编号 级别 位置 问题 后果
P0-1 高 module/dm/clean/CleanCache.java 全局静态人员集合,不分案件 跨案件串库/丢数据/CME
P0-2 高 module/dm/service/DmService.java:636-637 清洗完成删全局 TMP 根目录 其他用户在途临时文件蒸发
P0-3 高 common/cache/GlobalCache.java:64-127 静态 HashMap 运行期 clear+重建无同步 CME/半空缓存/模板错乱
P0-4 高 module/agent/service/impl/AgentChatServiceImpl.java:288-307 /chat/stream 无会话级互斥 同会话双流交错写记忆
P0-5 高 module/govern/controller/GovernController.java:56-65 治理任务无防重入+阶段异常被吞 并发治理互删数据/半治理落库
P0-6 高 module/aiclean/service/AiSmartCleanService.java:116,450 synchronizedList 被 computeIfAbsent 击穿 裸 ArrayList 并发写
P0-7 高 common/cache/CaseDataCache.java:68-94 三套互不互斥的锁保护同一结构 开案初始化期并发读 CME
P1-8 中 AgentChatServiceImpl.java:529-535 等 message_count 读改写遗留 计数漂移
P1-9 中 common/utils/GlobalPool.java 全局 16 线程池被长任务占满 其他用户任务无限排队
P1-10 中 SseService.java:133 等 4 个 @Scheduled 单线程调度器+回收与后台任务无互斥 心跳延迟/在跑任务夭折
P1-11 中 module/agent/service/AgentService.java:96,146-165 agent 池无上限无淘汰 内存泄漏
P1-12 中 common/utils/SystemConfManager.java:18,41-63 无锁 HashMap+并发写文件覆盖 配置丢失/读到 null
P1-13 中 module/agent/attachment/AttachmentDocIndexer.java:80 局部 Semaphore 无全局上限 嵌入模型被打爆
P1-14 中 DmService.java:672-679 先启后台清洗后删旧数据 写删交错
P1-15 中 module/agent/python/PythonExecutor.java:235-259 快照指纹 check-then-act 偶发执行失败
P1-16 中 AgentChatServiceImpl.java:494-500 等 接口缺案件归属一致性校验 跨案件读写会话
P1-17 中 GovernService.java:153、SystemService.java:189 GOVERN_CONFIG 全表 clear 误伤其他案件在用条目
P1-18 中 common/utils/LicenseUtil.java:44 static synchronized 锁内起子进程 全局串行瓶颈
P1-19 中 DmService.java:860-907 deletedFiles 多步删除无事务 孤儿数据
P2-20~29 低 见第三节 各项 影响小/低频

1. P0 高危(多用户并发必现:串数据 / 丢数据 / 崩溃)

P0-1 CleanCache 全局静态人员集合,跨案件串写 + 非线程安全

位置:ai-server/src/main/java/com/zsjz/ai/module/dm/clean/CleanCache.java:17-49

现状:

public static final HashMultimap<String, String> PERSON_LIB_NO_MAP = HashMultimap.create();
private static final Set<String> PERSON_LIST = Sets.newHashSet();
public static final Set<PersonBasicInfo> PERSON_LIST_A = Sets.newHashSet();
public static final Set<PersonLibNo> PERSON_LIB_NOS = Sets.newHashSet();
  • 写入方:清洗 Loader(AbstractCallRecordDataLoader.java:88、AbstractTransRecordDataLoader.java:55、CallOpenInfoDataLoader.java:63、TransOpenDataLoader.java:71)并发调 addPerson(:22-28,synchronized)与 addPersonLibNo(:43-45,无锁)、putPerLib(:47-49,无锁)。
  • 消费方:PersonLibNoService.flushPerLib()(PersonLibNoService.java:16-35)无锁迭代 PERSON_LIB_NOS(Iterables.partition)+ finally clear();PersonBasicInfoService.flushPerson()(PersonBasicInfoService.java:171-181)连续两次调 getPersonList()(CleanCache.java:30-41,不加锁地遍历 PERSON_LIST、补写 PERSON_LIST_A、PERSON_LIST.clear())。
  • 触发点:DmService.java:599 任一用户清洗完成时 CleanCache.PERSON_LIB_NO_MAP.clear() 全局清空。

多用户触发场景:用户 A(案件 1)与用户 B(案件 2)同时清洗。这些集合不区分案件:A 先完成清洗时,回调把全局缓冲里 B 案件累积的人员/手机号/银行卡也一并 insertOrUpdateBatch 进 A 的案件库,随后 clear() 把 B 尚未 flush 的条目清掉——B 自己的 flush 拿到空/缺数据。

后果:跨案件数据污染(A 库里出现 B 案件的人员)+ 人员库记录丢失 + 并发 add/iterate 非线程安全集合的 ConcurrentModificationException。这是全部发现中唯一会直接把数据写错库的问题,最高优先级。

修复方案(按案件分桶,参照 CaseDataCache 已有模式):

  1. 把 4 个静态集合收敛进 Map<Long, CleanBuffer> BUFFERS(CleanBuffer 内含这 4 个集合),key 为 caseId。清洗链路全部跑在 CaseContextHolder.runWith 内(DmService.java:669-672 已绑定),取 caseId 没有障碍。
  2. 同一案件内的并发 add 用 synchronized(buffer) 保护;addPersonLibNo/putPerLib/addPerson 增加 caseId 参数或从上下文解析。
  3. flushPerLib()/flushPerson() 改为只取当前案件的 buffer:BUFFERS.remove(caseId) 原子取出后落库(remove 后即无人再写,迭代无需加锁)。
  4. getPersonList() 的"补齐 + 清空"逻辑整体收进 synchronized(buffer)。
  5. 删除 DmService.java:599 的全局 PERSON_LIB_NO_MAP.clear();findPersonVal 查询改为查当前案件 bucket(CaseDataCache.findPersonVal 已有同款实现可对齐)。

注意:不能走"加个 synchronized 就行"的最小改动路线——问题不只是线程安全,跨案件串写必须靠分桶解决。


P0-2 清洗完成回调删除全局共享 TMP 根目录

位置:ai-server/src/main/java/com/zsjz/ai/module/dm/service/DmService.java:636-637

现状(清洗全部完成后的回调,跑在虚拟线程上):

//删除临时文件
log.info("删除临时文件");
FileUtil.del(PathConst.TMP_PATH);
FileUtil.mkdir(PathConst.TMP_PATH);

PathConst.TMP_PATH(QingJian/tmp)是全局唯一的临时根目录,所有用户的上传暂存、解析中间产物都在这里。同文件 :779,794 的 cancelClean/clearUploadBatch 已经是按 tmp/{batchId} 粒度管理的,唯独这里整树抹除。

多用户触发场景:用户 A 清洗完成触发整树删除时,用户 B 正在上传/预解析,其 tmp/{batchId} 下的文件随机消失。DmService.java:153-154 的注释(预处理产物因此特意不落 TMP_PATH)说明这个坑之前已经踩过一次。

后果:B 的解析/清洗中途文件缺失报错,且偶发、难排查;多人同时上传时是高概率事故。

修复方案:

  1. 回调签名里已有 batchId(DmService.java:590),把整树删除改为只删本批次目录 FileUtil.del(PathConst.TMP_PATH.resolve(String.valueOf(batchId)))。
  2. 兜底:新增一个低频定时任务(如每日)清扫"修改时间超过 24h 的批次目录",替代"用完即删根"的语义。孤儿批次本就该有 TTL(对齐 ChatAttachmentService.java:384 孤儿附件清理的做法)。

P0-3 GlobalCache 静态 HashMap 运行期 clear+重建,与全链路并发读无同步

位置:ai-server/src/main/java/com/zsjz/ai/common/cache/GlobalCache.java:64-127

现状:8 个 public static 容器全部是普通 HashMap/ArrayList(TEMPLATE_TABLE_MAP、ID_NAME_ZH_MAP、TABLE_FIELD_CACHE、TABLE_INFO_LIST、ALL_FIELDS_MAP、ALL_MATCHED_FIELDS_MAP、TABLE_HEAD(引用还是非 final)、DIRECTION_CONF_MAP)。initCache()(:111-116)先 cleanCache() 全清(:118-127)再查库逐个 putAll。

运行期写触发点(不是只在启动时):

  • DmService.java:596——任一用户清洗全部完成的回调里重建;
  • TableInfoService.java:67、:89——用户在平台上创建/更新模板时(事务内重建,见 P1 补充);
  • SystemService.java:184-190——/plat 的 cleanCache 接口,任何登录用户可调。

读触发点:清洗规则匹配(各 *DataCleaner)、PreDataListener.matchTemplate、TemplateCandidateLoader.java:34-35、GovernService.calcTreeTemplateDataCnt(:109 直接 stream 遍历 TABLE_INFO_LIST)、表头渲染、TimeSeriesService、FunRegular、AiCleanService 等 15+ 处主链路。

多用户触发场景:用户 A 清洗完成/建改模板触发 clear()+putAll() 的几百毫秒窗口内,用户 B 的清洗、AI 判定、治理正在迭代这些 Map。

后果:ConcurrentModificationException;或读到半空缓存导致模板匹配失败、清洗规则缺失、方向判定错乱;HashMap 扩容期并发读还可能无限循环(JDK7)/数据撕裂。

修复方案(copy-on-write 快照整体替换):

// 字段改 private static volatile,对外提供只读访问器;TABLE_HEAD 引用一并 volatile 化
private static volatile Map<Integer, TableInfo> templateTableMap = Map.of();
public static Map<Integer, TableInfo> templateTableMap() { return templateTableMap; }

public static synchronized void initCache() {
    // 1. 全部在局部变量里构建(查库、分组、装 Map/List),期间旧引用原样可读
    // 2. 构建完成后按依赖顺序一次性替换引用
    this.templateTableMap = newTemplateTableMap; ...
}

要点:

  1. 绝对不要在共享 Map 上先 clear 再 put——哪怕换成 ConcurrentHashMap 也只是不崩,仍然会读到半空数据;必须"构建新容器 → 整体替换引用"。
  2. 读方拿到引用后只读不写;TABLE_INFO_LIST 这类 List 换成不可变 List.copyOf(...),杜绝读方迭代时被改。
  3. TableInfoService.create/update 的重建移到事务提交后执行(TransactionSynchronization.afterCommit),避免事务回滚后缓存与库不一致(顺带修掉事务内重建的问题)。

P0-4 /chat/stream 无会话级互斥,同一会话可并发双流

位置:ai-server/src/main/java/com/zsjz/ai/module/agent/service/impl/AgentChatServiceImpl.java:288-307

现状:

public Flux<ServerSentEvent<String>> stream(ChatRequestDTO req) {
    ...
    return Flux.defer(() -> doStream(req)) ...   // 无任何锁

doStream(:316 起)内 agentService.getOrCreateAgent(...) 取共享的 HarnessAgent 实例(key 四维含 userId,但同用户的同 (专家×模型×工作区) 是同一个实例),随后 agent.streamEvents(msgs, ctx) 写记忆槽 (RUNTIME_USER_ID="default", sessionId) 与 AgentScope 工作区 transcript/events 文件。

对照:研判链路有完整双保险——应用锁 sessionLocks(InsightServiceImpl.java:136)+ DB 状态位 acquireStreamLock(InsightServiceImpl.java:307-333,doFinally 复位)。chat 链路两道都没有。

多用户触发场景:同一用户双标签页/双击重试/刷新重连,对同一 sessionId 同时发消息:两个流并发订阅,交错写同一记忆槽与同一工作区文件,且与 P1-8 的读改写叠加造成计数漂移。

后果:对话记忆与转录交错(刷新后历史与模型上下文互相污染)、工作区文件并发写异常、计数漂移。

修复方案(照抄研判链路的双保险):

// 1) 应用锁:ConcurrentHashMap<String, ReentrantLock> chatSessionLocks(key=sessionId)
// 2) DB 状态位:agent_chat_session 增加 streaming 状态列(或复用状态机字段),CAS 抢占
ReentrantLock lock = chatSessionLocks.computeIfAbsent(key, k -> new ReentrantLock());
if (!lock.tryLock(sessionLockTimeoutMs, TimeUnit.MILLISECONDS)) {
    return Flux.just(sse("error", Map.of("type", "error", "error", "该会话正在回复中,请稍候")));
}
try {
    if (sessionMapper.acquireStreamLock(sessionId) == 0) { /* 同上,返回错误帧 */ }
    return doStream(session, req).doFinally(sig -> {
        sessionMapper.releaseStreamLock(sessionId);
        lock.unlock();
    });
} catch (InterruptedException e) { ... }

实施时与 P2-21 一并处理:删除会话时不要无条件 sessionLocks.remove(研判链路 InsightServiceImpl.java:238 的同款问题),至少先 tryLock 探测或干脆只依赖 DB 状态位拒绝。


P0-5 治理任务无防重入;阶段异常被吞后流水线照常执行

位置:

  • ai-server/src/main/java/com/zsjz/ai/module/govern/controller/GovernController.java:56-65
  • ai-server/src/main/java/com/zsjz/ai/module/govern/serivce/GovernService.java:110,153,360-367

现状:

// GovernController.start() —— 直接提交,无任何"该案件是否在治理中"检查
ThreadUtil.execAsync(() -> governService.calcTask(
        Objects.requireNonNullElse(taskProcess, 66), caseId, userId));
  • GovernService.calcTreeTemplateDataCnt()(:108-115)第一步就是 truncate govern_tree,阶段任务还有 copyToCallRecord 等对共享表的全量重写(:245-248);
  • GovernService.java:153 CaseDataCache.GOVERN_CONFIG.clear() 全表清(见 P1-17);
  • awaitStage(:360-367)捕获 CompletionException 只记日志,executeFirstStageTasks 失败后第二/三阶段照常执行。

多用户触发场景:双击按钮、或同一用户两标签页对同一案件同时点"开始治理"→ 两条流水线并发,互相 truncate 对方刚写入的数据;单条流水线阶段一失败后阶段二/三拿着半成品数据继续跑,SSE 还在向用户报"正常"进度。

后果:治理树/清洗结果表数据互相覆盖、半治理状态落库、用户看到成功但数据是坏的。

修复方案:

  1. 防重入:参照 AiSmartCleanService.upsertBatch(状态位防重入,AiSmartCleanService.java:109-115)或 InsightServiceImpl 的"应用锁 + DB 状态位"两级防护:治理开始前 CAS 抢占状态位(可用现成的 ai_clean_batch 状态表或新增 govern_run 表),已在跑则直接返回"治理进行中"。
  2. 失败中断:awaitStage 捕获异常后应抛出终止流水线(或返回 boolean 由调用方短路),并通过 SSE 推送失败事件;禁止"吞异常继续下一阶段"。
  3. 阶段任务内部对 truncate/copyToCallRecord 的调用天然要求"同一案件同时只有一条流水线",状态位解决后此项即闭环。

P0-6 AiSmartCleanService.reviewItems:synchronizedList 工厂被先 put 的裸 ArrayList 击穿

位置:ai-server/src/main/java/com/zsjz/ai/module/aiclean/service/AiSmartCleanService.java:116,450-457

现状:

:116  reviewItems.put(batchId, new ArrayList<>());                    // 放进去的是"裸" ArrayList
:450  reviewItems.computeIfAbsent(batchId,
              k -> java.util.Collections.synchronizedList(new ArrayList<>()))   // key 已存在 → 工厂永不生效
              .add(new NeedReviewItem(...));
:457  .set(AiCleanBatch::getNeedReviewDetail, Json.toStr(reviewItems.get(batchId)));

判定阶段按 sheet 派发多个虚拟线程并发调 addReview(同批次内并发),实际并发写的是 116 行放入的裸 ArrayList;:457 的序列化还与并发 add 竞争。批次结束 :233 才 remove。

后果:元素丢失 / ArrayIndexOutOfBoundsException / ConcurrentModificationException / 撕裂 JSON。同批次内必现风险(不跨用户,但 AI 智能清洗是主打功能)。

修复方案(一行改动 + 一处快照):

:116  reviewItems.put(batchId, java.util.Collections.synchronizedList(new ArrayList<>()));

// addReview 序列化前做同步快照,避免边写边读:
List<NeedReviewItem> snapshot;
List<NeedReviewItem> list = reviewItems.get(batchId);
synchronized (list) { snapshot = new ArrayList<>(list); }
... .set(AiCleanBatch::getNeedReviewDetail, Json.toStr(snapshot));

P0-7 CaseDataCache:同一结构被三套互不互斥的锁保护

位置:ai-server/src/main/java/com/zsjz/ai/common/cache/CaseDataCache.java:68-94,135-146,167-170,178

现状(CaseBucket 内是 HashMultimap + ReentrantReadWriteLock):

  • initCache()(:68-77)用 synchronized(bucket) 做 clear——读写锁不感知 synchronized;
  • initPersonLibNoCache()(:79-94)向 personLibNoMap 批量 put 完全无锁;
  • 读侧 getPerPhones(:135-138)、getPerCards(:143-146)、containsPersonValue(:167-170)读 multimap 不加读锁(只有 findPersonVal:178 加了);
  • 写侧 putPerLib(:102,122)拿写锁。

四条路径对同一个 Guava HashMultimap(非线程安全)各用各的并发控制。

多用户触发场景:同案件内,用户 B 开案 initCache(clear+无锁批量 put)期间,用户 A(或同一用户的 AI 工具线程)并发调 getPerPhones/containsPersonValue 读 multimap。

后果:CME 或 multimap 内部 HashMap 结构损坏 → 查询报错/返回错误归属。

修复方案(统一锁域,DB 调用移出锁外):

public static void initCache(Long caseId) {
    CaseBucket bucket = ...;
    List<PersonLibNo> rows = CaseContextHolder.callWith(caseId, () -> mapper.selectList(...)); // 锁外查库
    // 锁外把 personLibNoMap 的替代内容构建成新 HashMultimap,然后:
    bucket.writeLock.lock();
    try {
        bucket.personLibNoMap = newMultimap;   // personLibNoMap 改为非 final 字段,整体替换
        // personPhone/personCard 同理整体 putAll(它们是 ConcurrentMap,可增量 put,但为一致性建议同窗口替换)
    } finally { bucket.writeLock.unlock(); }
}
// getPerPhones/getPerCards/containsPersonValue 全部补读锁(与 findPersonVal 同款)

要点:锁内绝不做 DB/IO;multimap 的"重建"走锁内整体替换引用而非逐条 put(与 P0-3 同思路)。


2. P1 中危(特定场景出错或性能塌陷)

P1-8 message_count 读改写遗留(4 处应改 SQL 原子自增/条件更新)

位置:AgentChatServiceImpl.java

  • persistAssistantMessage::510 先 selectById,:529-535 回写 session.getMessageCount() + 1——同文件其余 3 处落库已改为 setSql 原子自增,唯独这处遗留;
  • sendMessage::217-227 同型读改写;:229-239 首条消息并发时标题被两条请求各写一次;
  • deleteMessage::257-264 Math.max(0, count - 1);
  • toggleStarMessage::272-285 读 !entity.getStarred() 回写(togglePinSession:170-177 同型)。

触发/后果:多端并发(P0-4 场景)时计数越用越小(丢更新)、置顶/收藏状态回跳。实体无 @Version(AgentChatSession.java:38-39 已确认),@Transactional 救不了常量回写。

修复方案:

// 自增/自减一律 setSql(PG 语法):
.setSql("message_count = message_count + 1")
.setSql("message_count = GREATEST(0, message_count - 1)")
// 翻转一律 SQL 端做:
.setSql("starred = NOT starred")
// 首条标题:仅在标题为空时写,消除双写:
.set(AgentChatSession::getTitle, title)
.isNull(AgentChatSession::getTitle)   // 加进 update wrapper 的 where 条件

P1-9 GlobalPool:全项目唯一共享池,长任务占满 → 其他用户无限排队

位置:common/utils/GlobalPool.java:19,35-47。固定 min(CPU×2, 16) 线程、LinkedBlockingQueue(100000)、默认 Abort 拒绝策略(10 万积压后才触发)。治理任务单阶段 orTimeout 10-40 分钟(GovernService.java:344-358),一个大案件即可占满全池。

修复方案:治理这类长任务改走虚拟线程 + 按案件信号量(同案件同时只允许一条流水线,与 P0-5 的状态位二选一即可,信号量是进程内的第二道);GlobalPool 只留短任务。拒绝策略显式设为 Abort 并把队列水位(如 >80% 打 warn 日志)暴露出来,避免过载表现为"无限排队"。

P1-10 @Scheduled 单线程调度器;数据源回收与在跑后台任务无互斥

位置:SseService.java:133(5s 心跳,遍历所有 emitter)、CaseDataSourceReaper.java:34(5min)、AiCleanWatchdog.java:76(60s)、ChatAttachmentService.java:384(1h,逐个删文件+向量行)。

两个问题:

  1. Spring 默认单线程调度器(App.java 只有 @EnableScheduling,未配 pool-size):附件清理/看门狗变慢时,全服用户的 SSE 心跳整体延迟 → 空闲连接被代理断开。
  2. CaseDataSourceReaper 回收(及 CaseSessionListener.java:40-62 登出即释放)不看"该案件是否有在跑任务":用户 token 过期/登出瞬间,其仍在跑的治理(40 分钟级)会因案件数据源被 close 而中途夭折。

修复方案:

  1. application.yaml 增加 spring.task.scheduling.pool.size: 4(一行配置)。
  2. 回收前检查活跃批次:查 ai_clean_batch/治理状态位中该案件是否有 RUNNING 记录,有则本轮跳过(reaper 已是保守策略风格,加这条与之一致)。

P1-11 AgentService.agentPool:无上限、无空闲淘汰,computeIfAbsent 内做 DB 查询

位置:module/agent/service/AgentService.java:96,146-165,373-384。池 key 四维 w{ws}-a{agent}-m{model}-u{user},只有删除专家时 evict;换模型/换案件会不断累积重量级 HarnessAgent 实例。对照 InsightAgentFactory(POOL_MAX=64 + 30 分钟空闲回收)已有正确范式。另 computeIfAbsent 的 bin 锁内做两次 DB 查询,同 bin 其他 key 被阻塞。

修复方案:对齐 InsightAgentFactory——容量上限 + lastAccess 时间戳 + 惰性/定时空闲回收;DB 查询移到 computeIfAbsent 外(先查配置再进 computeIfAbsent,双检即可)。

P1-12 SystemConfManager:无锁 HashMap + 并发写文件互相覆盖

位置:common/utils/SystemConfManager.java:18,41-63。普通 HashMap 无锁读写;set/setInitFlag/setForceDataFlag 都是"put + 全量重写文件",两个写并发时后写文件者覆盖前者的全部键;扩容期 get 可能短暂 null。

修复方案:systemConf 换 ConcurrentHashMap;写方法加 synchronized(读写锁更佳但此处写极少,synchronized 足够);写文件改为"写临时文件 → 原子 move 替换",防止写一半被读到。

P1-13 AttachmentDocIndexer:Semaphore 是方法局部变量,无全局嵌入并发上限

位置:module/agent/attachment/AttachmentDocIndexer.java:80。闸门只限"单个附件内"的 chunk 并发;多用户同时上传大附件时,N 个附件 × indexConcurrency 路并发同时打向全局共享的嵌入模型。正确范式在 FileRecognitionService.java:99-126(成员级公平 Semaphore + maxPending 积压上限)。

修复方案:把 Semaphore 提为 Bean 成员(全局单例,permits = 全局嵌入并发预算),可另加全局 inflight 上限排队(对齐 FileRecognitionService)。

P1-14 doClean:先启动后台清洗,请求线程随即删旧数据

位置:DmService.java:672-679。:672 先 Thread.startVirtualThread(...executeCleanJobs...),:674-678 紧接着 deletedFiles(...)(reClean=1 时)——新清洗任务写库与"清理旧数据"(删业务表数据+Lucene)并发交错,存在删掉新写数据的窗口。

修复方案:把 deletedFiles 移到启动后台线程之前(请求线程先删后启);或作为后台任务的第一步在清洗线程内执行。一行顺序调整。

P1-15 PythonExecutor 快照指纹 check-then-act;Windows 回退路径打不开

位置:module/agent/python/PythonExecutor.java:235-259(读 meta → 比对指纹 → exportSnapshot → 写 meta,全程无锁);DuckdbSnapshot.java:100-109(目标被 Python 子进程占用时替换失败 → 回退原库路径,而原库被 Java 进程写锁持有,子进程打不开)。另 exportSnapshot 未 synchronized,与 CaseDataSourceRegistry.close(CaseDataSourceRegistry.java:178-213)存在窗口。

修复方案:按案件(或源库路径)加 per-key 锁把"读指纹→导出→写指纹"串起来;替换失败不要回退原库路径,改为重试一次或明确报错。低频路径,可与 P2-23(LRU 误删产物)一并处理。

P1-16 chat/insight 接口缺"会话属于当前打开案件"的一致性校验

位置:AgentChatServiceImpl.requireSession:494-500(只查存在性);InsightServiceImpl.requireSession:531-540(校验了 session.caseId == 入参 caseId,但入参 caseId 本身未校验是否是当前用户打开的案件)。AgentChatServiceImpl.java:491 注释表明"项目已移除用户体系、按 caseId 归属"是有意设计,但目前连"会话的 workspace == 当前打开案件"这一层也没校验,与 CaseInfoService.assertOwner:260-265、附件三重校验的口径不一致。

后果:拿到别处泄漏的会话 ID(日志/导出/分享链接)即可跨案件读写消息、删除会话。不属严格意义的"并发竞态",但属于多用户数据隔离面,一并列出。

修复方案(最小改法,不动"无用户体系"的设计):requireSession 增加与 CaseContextHolder.currentCaseId() 的一致性校验——读接口可放宽,写操作(stream/delete/pin/star)必须一致。

P1-17 GOVERN_CONFIG 按 key 已分案件,但 clear() 是全表清

位置:GovernService.java:153、SystemService.java:189。key 形如 cache_getCallTableName{caseId}(StrConsts.java:12-18),天然按案件隔离;但治理启动/清缓存时 clear() 是全表清:A 用户治理进行中,B 用户开治理会清掉 A 正在用的条目(后续 computeIfAbsent 按当时 DB 状态重算,可能中途换表)。

修复方案:按 key 前缀清除:GOVERN_CONFIG.keySet().removeIf(k -> k.startsWith(prefixOf(caseId)));SystemService.cleanCache 全清保留(那是管理端显式清缓存语义)。容器本身是 ConcurrentMap,removeIf 安全。

P1-18 LicenseUtil:static synchronized 全局锁内起 OS 子进程

位置:common/utils/LicenseUtil.java:44(static synchronized checkAuthStatus → getAuthCode() → ProcessBuilder.start() + waitFor)。所有线程的授权校验全局串行,每次校验一次进程创建。当前 LicenseFilter 是透传(LicenseFilter.java:39-42),仅个别入口触发,影响有限;一旦恢复全局拦截就是全站串行瓶颈,故列 P1 备忘。

修复方案:校验结果加 TTL 缓存(如 60s 内直接返回上次结果);或把子进程调用移出锁外(锁只保护结果字段的读写)。

P1-19 deletedFiles 多步删除无事务

位置:DmService.java:860-907。先删 file_info 行,再并发删各业务表数据 + Lucene(ES)数据,任一步失败留下孤儿业务数据(残留可查询的脏数据)。与 P1-14 的时序问题叠加放大。

修复方案:DB 删除(file_info + 业务表)包进一个 @Transactional;ES 删除移到事务提交后(afterCommit)执行并允许重试/记录补偿日志(ES 侧失败不应回滚 DB)。


3. P2 低危(记录在案,低频/影响小,择机处理)

编号 位置 问题 后果
P2-20 SseService.java:54-58 vs 119-128 closeSee 两步 remove+complete 与 see() 重连(put+old.complete)竞态:remove 可能 complete 掉刚重连的新 emitter 同用户"退出+重连"并发时前端掉线,补发后不恢复。修法:closeSee 用 emitters.remove(userId, emitter) 两参版本(只删自己拿到的那个引用)
P2-21 InsightServiceImpl.java:136,238,328-333 sessionLocks 仅删会话时清理(常驻缓慢泄漏);删会话时无条件 remove 可能移走正被持有的锁 锁失效窗口 + 内存缓慢累积。修法:不 remove(条目极小),或 remove 前 tryLock 探测;与 P0-4 一并改
P2-22 AgentService.java:373-384 专家编辑后 evict 实例,在途流继续用旧配置至本轮结束 旧提示词/工具组再生效一轮。可接受则文档注明即可
P2-23 PythonExecutor.java:313-325 evictOldExecs LRU 按目录时间删除,同案件并发执行时 listing 与删除间无协调 超额时刚产出的图片被另一并发执行的清理删掉(前端 404)。修法:只清理"非活跃 execId"目录
P2-24 UsageStore.java:41-49 超限时并发 removeFirst 会多删事件;COWList removeFirst 是 O(n) 监控统计轻微失真。修法:synchronized 包住超限检查+删除即可
P2-25 InsightServiceImpl.java:171-175 createSession 的 countByCase >= max 是 check-then-act 并发创建可超上限(软限制,影响小)
P2-26 AgentChatServiceImpl.java:65,332、InsightServiceImpl.java:98,363 全员共用 RUNTIME_USER_ID="default" 当前靠 UUID 派生 sessionId 兜底不串;但若未来 sessionId 复用/跨环境合并 stateStore 即串记忆。架构性隐患,记录即可
P2-27 TowerUtils.java:47-52 init() 的 FLAG.get() check-then-act 非 CAS 并发首次触发重复 cell_init(native 重复加载)。修法:FLAG.compareAndSet(false,true) 进入后再加载
P2-28 DuckdbUnpooledDataSource.java:142-145 setNetworkTimeout(Executors.newSingleThreadExecutor(), ...) 每连接新建 executor 且从不关闭(JDBC 规范由调用方负责关闭) 当前 defaultNetworkTimeout 未配置,属死代码级隐患;一旦启用即每连接泄漏一个线程池。修法:不设 networkTimeout 或复用共享 executor
P2-29 McpCaseSession.java:57 sessionCases MCP 会话结束无回调清理 条目极小,长期运行缓慢累积。修法:MCP 会话关闭钩子里 remove

4. 修复路线图

第一批(小改动、独立、当天可完成,先消掉"必现"类)

项 改动量
P0-2 回调改删 tmp/{batchId} + 过期清扫定时任务
P0-6 一行 put 改 synchronizedList + addReview 快照
P1-8 4 处 setSql/条件更新
P1-14 调整删除与启动的顺序
P2-20 remove(key, value) 两参版本
P2-27 compareAndSet
P1-17 removeIf 前缀清除

第二批(核心重构,1-2 天,建议每项独立提交 + 单测)

  • P0-1 CleanCache 按 caseId 分桶(连带删除 DmService.java:599 全局 clear、改造 flushPerLib/flushPerson 签名)
  • P0-7 CaseDataCache 锁域统一(DB 移出锁外 + 重建走锁内整体替换)
  • P0-3 GlobalCache 快照替换 + TableInfoService 重建移到 afterCommit
  • P0-5 + P1-9 治理防重入(状态位)+ awaitStage 失败中断 + 长任务脱离 GlobalPool

第三批(隔离与资源治理)

  • P0-4 chat 会话双保险锁(与 P2-21 一并)
  • P1-10 调度池扩容 + reaper 回收前查活跃批次
  • P1-11 agent 池容量上限 + 空闲回收
  • P1-13 全局嵌入信号量、P1-15 per-key 锁、P1-16 归属一致性校验

第四批(择机):P1-12、P1-18、P1-19、P2-22~29。


5. 并发回归验证建议

  1. 单测(可离线跑):用 CountDownLatch 齐发 N 线程的用例覆盖:
    • CleanCache:N 案件并发 add + 各自 flush,断言互不串、不丢(P0-1);
    • addReview 并发 1 万次 add,断言 size 精确(P0-6);
    • message_count 并发自增 1000 次,断言最终值 = 1000(P1-8);
    • GlobalCache 重建期间持续读,断言永不为空/不抛 CME(P0-3)。 参考现有 mvn -o -pl ai-server test -Dtest='com.zsjz.ai.**'(离线)。
  2. 双浏览器手工场景(真实多用户):两账号同时清洗不同案件(观察人员库归属、tmp 目录、进度 SSE 不串);同一会话双标签页同时发消息(第二路应收到"正在回复中"错误帧);双击"开始治理"(第二路应被状态位拒绝)。
  3. 观测哨:日志关键字 ConcurrentModificationException、ArrayIndexOutOfBounds、缺少案件上下文;ai_clean_batch 的计数列与实际行数对账。

6. 附录 A:已排查为安全的部分(勿重复排查)

领域 结论
上下文传播 CaseContextHolder 用 JDK ScopedValue + CaseContextFilter 请求作用域 + ReactorCaseContextConfig(onScheduleHook)+ 各异步链路显式 runWith/callWith,结构上消除 ThreadLocal 串案。残留风险仅为"新增异步提交忘绑定"(读路径静默降级打 error 日志,写路径抛异常,可发现)
附件链路 草稿附件 CAS 绑定(ChatAttachmentService.java:233-247)、孤儿条件删除(:396-402)——最近一轮已修,本轮复核无新问题
研判流互斥 会话锁 + DB 状态位 + doFinally 复位(InsightServiceImpl.java:307-333),双保险正确(仅 P2-21 的锁 map 清理时机小瑕疵)
案件数据源生命周期 CaseDataSourceRegistry synchronized 开关 + 配额、CaseRoutingDataSource 按案件物理 key 路由、DuckdbSnapshot 唯一别名 + 原子替换(:62-83)
SSE 连接管理 按 userId 的 ConcurrentHashMap + onCompletion/onTimeout/onError 三回调清理 + 心跳失败清理,无 emitter 泄漏(仅 P2-20 的两参 remove 竞态)
Python 沙箱 每次执行独立 tmp/py-{uuid} 目录,产物按 execId 隔离
文件名安全 产物/附件雪花 ID/UUID 前缀 + 路径穿越校验(ArtifactService.java:310-327、AgentPythonFileController.assertSafeSegment:85-99),多用户同名文件不互相覆盖
懒加载 DCL InsightPromptBuilder:79-96、AgentScopeMcpToolProvider:151-163、PythonExecutor:98-110、SearchDocIndexService.ensureIndex:59-74、OnnxEmbeddingEngine、StateManager(holder)——全部正确
全文检索 已是 Elasticsearch(非嵌入式 Lucene):SearchDocIndexService 按 caseId 分桶攒批 + per-buffer synchronized + _bulk 提交,caseId 强制打进文档与过滤条件,无跨案件风险。仅 flush 与并发 add 存在"最后一批 refresh 前不可搜"的毫秒级窗口,可接受
Redis 使用面极小(Sa-Token 会话 + SystemService/CacheManageService 的 keys(prefix)+delete 管理端清缓存)。keys() 会阻塞 Redis 但仅管理操作触发,可接受;无业务数据写 Redis 的 check-then-act
事务边界 6 个 @Transactional 类均为短 DB 事务;InsightServiceImpl 在 Reactor 线程正确改用 TransactionTemplate;未发现"事务内提交异步任务读未提交数据"模式。唯一边界:TableInfoService.create/update 事务内调 GlobalCache.initCache()(已并入 P0-3 修复)
线程卫生 CompletableFuture 均显式传 executor(无 commonPool 滥用);全项目未用 wait/notify、CountDownLatch、@Async;SSE/HTTP 客户端无每请求新建模式
单例 Bean 状态 全量扫描:单例 Bean 无请求期写入的可变实例字段(仅 @Value/@ConfigurationProperties);清洗策略经 CleanFactory 每次 supplier.get() 新建,非单例
Agent 实例隔离 AgentService.agentPool key 含 userId,用户间实例隔离正确(仅 P1-11 的容量/淘汰问题)

7. 附录 B:项目内可复用的并发防护范式(修 P0/P1 时对齐这些)

范式 参照实现 适用问题
案件+用户上下文绑定 CaseContextHolder.runWith/callWith(ScopedValue) 一切异步提交
按 caseId 分桶缓存 CaseDataCache.BUCKETS P0-1 CleanCache 分桶
快照整体替换缓存 GlobalCache.initCache()(已实施,见第九节) P0-3、P0-7 的重建窗口
两级互斥(应用锁+DB 状态位) InsightServiceImpl.java:307-333 P0-4、P0-5
批次状态位防重入 + 看门狗 AiSmartCleanService.upsertBatch + AiCleanWatchdog P0-5
成员级公平信号量 + 积压上限 FileRecognitionService.java:99-126 P1-13、P1-9
SQL 原子自增/条件更新 AgentChatServiceImpl 已改的 3 处 setSql P1-8
有淘汰的实例池 InsightAgentFactory(POOL_MAX + 空闲回收) P1-11
短事务 + afterCommit 外部 IO InsightServiceImpl.persistMessage(TransactionTemplate) P1-19、P0-3 的 afterCommit

9. 修复实施记录(2026-09-30)

9.1 已修复(本轮全部完成并测试通过)

编号 修复内容 关键改动
P0-1 CleanCache 按 caseId 分桶 CleanCache 重写:4 个全局静态集合 → Map<Long, CleanBuffer>,同案件 synchronized(buffer);新增 drain() 原子取走本案件快照;flushPerLib/flushPerson 改收快照参数;删除 DmService 回调里的全局 PERSON_LIB_NO_MAP.clear()
P0-2 tmp 只删本批次 + TTL 兜底 回调改删 tmp/{batchId};DmService 新增 sweepStaleTmpBatches()(每 6h 清 24h 前的批次目录)
P0-3 GlobalCache 快照整体替换 8 个容器改 private volatile + 只读访问器(17 个消费方已迁移);initCache() synchronized + 锁外构建 + 一次性替换;TableInfoService 重建移到 afterCommit;删除 cleanCache()
P0-4 chat 会话级互斥 AgentChatServiceImpl 新增 streamLocks(非阻塞 tryLock,抢不到立即错误帧),doFinally/catch 释放;同一会话双流被拒
P0-5 治理防重入 + 失败中断 RUNNING_CASES 原子占位(Controller tryStartGovern,calcTask finally 释放);awaitStage 改为上抛终止流水线;阶段二/开户信息链路 exceptionally 同步改上抛;失败经 SSE 推送「治理失败」
P0-6 reviewItems 同步容器 :116 放入 synchronizedList(工厂真正生效);addReview 序列化前锁内快照拷贝
P0-7 CaseDataCache 锁域统一 initCache 锁外查库构建新 multimap,写锁内整体替换引用(字段改非 final);getPerPhones/getPerCards/containsPersonValue 补读锁并锁内拷贝
P1-8 计数/状态原子化 sendMessage/persistAssistantMessage → setSql("message_count = COALESCE(message_count,0)+1");deleteMessage → GREATEST(...-1, 0);togglePin/toggleStar → setSql("pinned = NOT pinned");首条标题改条件写(title IS NULL OR title='新对话'),双请求仅一条生效
P1-9 治理脱离 GlobalPool 新增 GOVERN_EXECUTOR(虚拟线程每任务一条),3 处提交点替换;长任务不再挤占全服 16 线程池
P1-14 doClean 先删后启 deletedFiles 移到启动清洗线程之前
P1-17 GOVERN_CONFIG 按案件清除 新增 CaseDataCache.cleanGovernConfig(caseId)(前缀+剩余部分精确相等,规避 caseId=12 误伤 123),治理启动改用;SystemService 管理端全清保留
P2-20 SSE closeSee 竞态 改两参 emitters.remove(userId, emitter),不再误踢重连的新连接
P2-27 TowerUtils CAS init() 改 FLAG.compareAndSet,失败路径复位 FLAG 保留重试语义

第二批补记:P1-11/12/13/15/16/18/19 与 P2-21/23/24/25/28/29 已于同日补齐,见 9.4 节。仍明确不做的只有两类:P0-4 的 DB 状态位(多实例保险,单实例部署假设下进程内闸门已覆盖,多实例部署前须补)、P2-22(evictAgentPool 后在途流用旧配置跑完当轮——可接受的软失效,删除专家后旧实例被驱逐、新对话即用新配置)与 P2-26(RUNTIME_USER_ID="default" 架构性备忘,当前靠 UUID 派生 sessionId 兜底不串)。

9.2 测试结果

  • 新增并发回归测试 4 类 19 个用例,覆盖:多案件并发写入隔离与 drain 精确性、add/drain 竞争无 CME、putPerLib 读写锁并发、cleanGovernConfig 精确匹配(12 vs 123 陷阱)、GlobalCache 快照撕裂哨兵、治理占位并发唯一胜者——连跑 4 轮全绿。
  • 全量测试:356 run / 8 failed / 15 skipped。8 个失败(工具注册数断言 67 vs 55 等)已用 git stash 在未改动 HEAD 上原样复现,属上一轮功能提交遗留的过时断言,与本次修复无关。
  • 编译:主源码 + 测试源码离线编译零错误。

9.3 遗留的过时测试(非本次引入,建议另行修复)

AgentToolRegistryTest(5 处)、AgentCapabilityToolkitTest(1 处)、AgentScopeMcpToolProviderTest(2 处)断言的业务工具总数(67/55、24/12)与当前代码注册的实际工具数不一致,疑似某次工具组调整后未同步更新断言。

9.4 第二批修复实施(2026-09-30 同日补齐:P1-10/11/12/13/15/16/18/19 + P2-21/23/24/25/28/29)

编号 修复内容 关键改动
P1-11 AgentService 池上限 + LRU 淘汰 新增通用 AgentInstancePool<T>(容量 64 + 空闲 30 分钟,对齐 InsightAgentFactory;LRU 用单调递增序号而非墙上时间戳——同毫秒访问会让"最久未访问"判定退化,测试实测踩中);AgentService 的 getOrCreateAgent 改为「锁内快查 → 锁外查库构建 → 锁内 put 复用先到者」,DB/IO 不再进 computeIfAbsent 的 bin 锁
P1-12 SystemConfManager 并发安全 HashMap→ConcurrentHashMap;写方法 synchronized;落盘改「写 .tmp → 原子 move 替换」,进程任意时刻读到的都是完整 JSON
P1-13 附件嵌入全局闸门 AttachmentDocIndexer 的局部 Semaphore 提为成员级懒初始化单例(全部附件共享一份 permits),多用户同时上传时嵌入模型总并发受控
P1-15 Python 快照导出竞态 按 workspaceId 互斥整段「读指纹→比对→导出→写指纹」;删除"导出失败回退原库"(原库被 Java 写锁持有、Python 打不开,静默回退只会把问题推迟成难排查的执行失败),改为重试一次后显式抛错
P1-16 会话归属一致性校验 chat 写路径(发消息/删会话/删消息/置顶/收藏/改标题/流式)加 requireSessionInOpenCase:会话的 caseId ≠ 当前打开案件直接 403;insight 的 requireSession 同口径校验入参 caseId。读路径保持放宽(案件上下文取不到时放行,兼容无登录态调用)
P1-10 调度池扩容 + 回收互斥 spring.task.scheduling.pool.size: 4(心跳/看门狗/清理互相隔离);新增 CaseActivityProbe SPI——GovernActivityProbe(治理 RUNNING_CASES)与 AiCleanActivityProbe(ai_clean_batch 非终态批次,口径与看门狗一致)注册在案;Reaper 与 CaseSessionListener(登出/被踢)释放数据源前先探针,有在途任务则跳过本轮、交 Reaper 稍后兜底
P1-18 LicenseUtil TTL 缓存 checkAuthStatus 结果缓存 60s,多数调用不再进全局锁起子进程(原 static synchronized + ProcessBuilder 是全站串行瓶颈隐患)
P1-19 deletedFiles 一致性 跨主库/DuckDB 双数据源无法用单一事务,改用顺序性修复:先删业务表(任一表失败收集后整体抛出,file_info 保留作为重试入口),删净后再删 file_info 与 AI 档案,ES 尾删尽力而为。失败重试幂等(按 fileId 删)
P2-21 研判锁被误移 deleteSession 只在「未被持有」时移除锁(tryLock 探测),流在跑时保留条目(条目极小,不破坏互斥优先)
P2-23 py-out LRU 误删 ACTIVE_EXEC_IDS 记录在途执行,evictOldExecs 跳过在途目录,刚产出的图片不再被并发清理删掉
P2-24 UsageStore 多删 超限修剪改 synchronized + while(并发超限各自 removeFirst 会多删事件)
P2-25 会话创建超限 count 检查与 insert 收进按案件锁(原 check-then-insert 间隔几百毫秒构建期,双开必超限)
P2-28 networkTimeout executor 每连接 newSingleThreadExecutor(泄漏线程池)→ 守护线程共享单例
P2-29 MCP 会话槽泄漏 sessionCases 改有界 FIFO(256),淘汰时连带关闭被逐会话的案件数据源

第二批测试:新增 AgentInstancePoolTest(5 用例:LRU/容量/get 刷新 recency/标记驱逐/并发安全)、UsageStoreTest(并发修剪精确到上限)、McpCaseSessionSlotsTest(256 槽 FIFO 淘汰 + 数据源连带关闭,Mockito);AgentExpertTest 适配新池类型。全量 363 测试仅剩 8 个 HEAD 既有过时断言(工具注册数,与本轮无关);全部并发测试(26 个)连跑 3 轮全绿。

明确不做:P2-22(在途流用旧专家配置跑完当轮——软失效可接受)、P2-26(RUNTIME_USER_ID 架构备忘)、P0-4 的 DB 状态位(多实例部署前必须补,单实例进程内闸门已覆盖)。