Commit 31ed4846 by luoqi

perf(plan): 写入段改端到端分批 —— 取数早就分块了,产物却全程不释放

2026-09-04 生产 OOM 宕机 2.5 小时。内核日志:
  Killed process (node) anon-rss:6,321,668kB(6.03GB) total-vm:16.7GB
机器 14GB。进程被杀后整机 swap 抖死,sshd 与 Web 全部无响应,靠人工重启才能恢复。
docker-compose.prod.yml 那段注释 2026-08-26 就写过「这是止血不是根治:池子还在涨,
8G 迟早也会到顶」—— 那天到了。

## 病根不是"没分块",是"不释放"

取数早就是 CHUNK=2000 一块(prefetchForBatch 内),但**产物累积**在跨全量存活的 Map 里;
runPool 的闭包又把 latest/persona/visit/logRows 全部 context-allocate,
V8 要到 runAllForHost 整个 frame 结束才可能回收 → 第 3 步扫全池时它们全还在,那一下就是峰值。

堆账(550,049 命中患者 / 约 96 万条 hit,按 V8 对象布局建模):
  latestByPatient 1.93GB + hitsByPatient 1.06GB + persona/visit 0.14GB
  + activePlans 0.09GB + logRows 0.07GB ≈ 活跃保留 3.3GB,实测 RSS 6.03GB(≈1.8×)

## 改动

· prefetchForBatch 拆成两半:
  - prefetchSnoozeAnchors(scope, now) —— **整轮一次**。它的代价跟"终态且冷静期未到期的
    计划数"走、与患者数无关(生产实测全库 106 行)。拆出来是让这条纪律在类型上也成立:
    它拿不到 patientIds,想塞进批循环也塞不了。
  - prefetchOneBatch(scope, ids) —— 每批一次,删掉内部的 CHUNK 循环,切批权交给调用方。
· 写入段改成端到端批循环:每批 预取 → runPool 写库 → createMany 日志 → 释放。
  hit 的所有权交给本批局部数组后立刻 hitsByPatient.delete(pid),
  批末连同 latest/persona/visit/logRows 一起出作用域回收。
· 新增 resolvePlanBatchSize():默认 2000,夹在 [1, 20000](上限守 PG 32767 bind ——
  latest/persona 走 patientId in ids),**显式传 0 = single-shot 回退开关**
  (照 cold-import 的 resolveCohortBatchSize 同一形状)。
· 每 5 批 + 末批打一行 rss= / heapUsed= / heapTotal=(rss 沿用 cold-import 字样便于并排 grep;
  判 8G 堆余量只能看 heapUsed,rss 含 Prisma Rust 引擎与碎片)。

## 唯一的硬全局依赖,以及它的坑

第 3 步 stale-close 拿本轮命中集对整个召回池取补集(`!hitsByPatient.has(pid)` → supersede)。
分批后 hitsByPatient 被逐批清空 ⇒ 改用跨批累加的 `hitPatientIds` Set(只存 id,约 54MB)。
 第 3 步**必须等所有批跑完再跑一次**。搬进循环里 = 每批把其他批的患者判成"信号消失"
   → 整池 supersede + 每条认领单一条 auto_release,**而且不报错**,
   那行「本轮关闭 N 条无信号 plan」的日志反而会解释成「这不是 bug」。
 循环体里**不加 try/catch**:今天的语义是"selectHits 抛错发生在任何写之前 → 整轮中止、
   第 3 步不跑";分批后前 N 批已落库,吞错"把剩下的批跑完"会让第 3 步拿着残缺命中集去关池子。

## 测试

新增 describe.each 覆盖 PAC_PLAN_BATCH_SIZE ∈ {0,1,2,3,2000},同一 fixture 各跑一遍,
断言引擎返回八个字段 + plan 终态 + 日志行**逐字段一致**,直接编码「结果与批大小无关」。
**为什么必须新加**:既有 20+ 条用例每个只有 1~4 个患者、默认批 2000,永远只有一批 ——
分批写错它们全绿。batch=1 是最凶的一档,同时压测跨批不误关 / touchedPlanIds 按批等价 /
snooze 外提后仍被每批查到。fixture 里特意放了一个本轮 0 命中的 p9 压 stale-close。

有牙验证:把第 3 步判据改回 hitsByPatient.has(分批后已清空)→ plansClosed 0→3、
active 与 assigned 双双变 superseded,多条用例立刻红。

## 预计效果与未做的部分

峰值从"latest+hits+logRows 同时在堆"的堆叠拆成单批驻留,预计 3.3GB → 约 2.1GB。
️ 这一级**不碰任何 SQL、不碰 selectHits**,08-30 那根 ×2.31 的并发杠杆与 work_mem 收益
   一个字节不动。但 hitsByPatient 仍在第 1 步满载(1.06GB),峰值仍与命中总数线性相关。
   要真正解耦需要二级(患者域分区 selectHits),那要改 SQL 且必须先在生产标定
   11×P 次子场景查询的墙钟 —— 生产当前宕机,标不了,另开。
️ 第 3 步的 activePlans findMany 仍是无分页全池扫(约 94MB,占比 3%)。改键集分页有真实风险:
   plan-engine-stale-close.spec 的 findMany mock 忽略 take/orderBy/id.gt,分页写错在 CI 里全绿。
   单独一件事做,别和本次混。
️ PAC_PLAN_BATCH_CONCURRENCY 从来没被 withCohortDerivedPool 算进连接池(2026-09-02 已因
   漏算一次导致健康日报连续三天超时)。本次不改,记在这。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
parent 3e86060c
import { PlanEngineService } from '../src/modules/plan/engine/plan-engine.service';
import { PlanEngineService, resolvePlanBatchSize } from '../src/modules/plan/engine/plan-engine.service';
import type { ScenarioHit } from '../src/modules/plan/engine/scenario.interface';
/**
......@@ -1156,3 +1156,109 @@ describe('归因继承的边界 — 判据是「客服碰过没」', () => {
for (const [args] of called) expect(args.where.planId.in).not.toContain('p-g');
});
});
/**
* 🔴 2026-09-04 端到端分批改造的**核心判据**:结果与批大小无关。
*
* 事故背景:生产 pac-service 涨到 6.03GB 常驻被内核 OOM killer 打掉,整机 swap 抖死 2.5 小时。
* 堆账里 latestByPatient(1.93GB)+ hitsByPatient(1.06GB)是大头,而它们原先**全程不释放**
* —— 取数早就是 2000 一块,但产物累积在跨全量存活的 Map 里,runPool 的闭包又把它们
* context-allocate 到整个 runAllForHost frame 结束。改造把取数分块升格成端到端分批。
*
* 为什么这组用例是必须的:上面那 20+ 条既有用例**一条都跑不到分批路径的差异** ——
* 它们每个只有 1~4 个患者,默认批大小 2000,永远只有一批。分批写错了它们全绿。
* 这里用 describe.each 把同一份 fixture 在不同批大小下各跑一遍,直接编码
* 「结果与批大小无关」这条不变式。
*
* ⚠️ batchSize=1 是最凶的一档:每批一个患者,同时压测三件事 ——
* ① 跨批不误关(stale-close 拿的是跨批累加的 hitPatientIds,不是被逐批清空的 hitsByPatient)
* ② touchedPlanIds 按批算与全量算等价
* ③ snooze 提在循环外之后仍能被每一批查到
*/
describe('端到端分批 — 结果必须与 PAC_PLAN_BATCH_SIZE 无关', () => {
const prev = process.env.PAC_PLAN_BATCH_SIZE;
afterEach(() => {
if (prev === undefined) delete process.env.PAC_PLAN_BATCH_SIZE;
else process.env.PAC_PLAN_BATCH_SIZE = prev;
});
/// 与上面「混合一批」同一份 fixture:5 种结局各一,外加一个 0 命中的患者压 stale-close
const runMixed = async () => {
const { prisma, plans, logs } = makeStore({
plans: [
{ id: 'u-old', patientId: 'p2', status: 'active', priorityScore: 50, reasons: [{ scenario: SCEN, subKey: 'k@1' }] },
{ id: 's-old', patientId: 'p3', status: 'active', reasons: [{ scenario: SCEN, subKey: 'k@old' }] },
{
id: 't-term',
patientId: 'p4',
status: 'abandoned',
snoozedUntil: new Date('2026-12-01T00:00:00Z'),
reasons: [{ scenario: SCEN, subKey: 'k@4' }],
},
// ⭐ p9 本轮 0 命中且有 active plan → 必须被 stale-close 关掉。
// 分批写错(第 3 步进了循环 / 用了被清空的 hitsByPatient)时,
// 这里会连 p1~p4 一起关掉,plansClosed 从 1 变 5 —— 这条就是照妖镜。
{ id: 'stale-old', patientId: 'p9', status: 'active', reasons: [{ scenario: SCEN, subKey: 'k@gone' }] },
],
});
const res = await engine(
prisma,
makeScenario([
hit('p1', 'k@1'),
hit('p2', 'k@1', 50),
hit('p3', 'k@new'),
hit('p4', 'k@4'),
]),
).runAllForHost({ hostId: HOST, tenantId: TENANT, now: NOW });
return { res, plans, logs };
};
/// 把最终状态压成一个可比对的快照(引擎返回值 + 落库结果),批大小不该改变它任何一位
const snapshot = (r: Awaited<ReturnType<typeof runMixed>>) => ({
result: {
scenariosRun: r.res.scenariosRun,
patientsHit: r.res.patientsHit,
plansCreated: r.res.plansCreated,
plansSuperseded: r.res.plansSuperseded,
plansUnchanged: r.res.plansUnchanged,
plansSuppressed: r.res.plansSuppressed,
plansClosed: r.res.plansClosed,
plansSkippedAssigned: r.res.plansSkippedAssigned,
},
// plan 终态:按 (patientId, version) 排序,避免并发写入顺序带来的假差异
plans: r.plans
.map((p) => `${p.patientId}|v${p.version}|${p.status}|${p.priorityScore}`)
.sort(),
logs: r.logs.map((l) => `${l.patientId}|${l.status}`).sort(),
});
const SIZES = [0, 1, 2, 3, 2000];
let baseline: ReturnType<typeof snapshot> | null = null;
test.each(SIZES)('批大小 %s → 结果与基线逐字段一致', async (size) => {
process.env.PAC_PLAN_BATCH_SIZE = String(size);
const snap = snapshot(await runMixed());
// 先自证这份 fixture 真的走到了 5 种结局,否则"全相等"可能只是都没跑
expect(snap.result).toMatchObject({
patientsHit: 4,
plansCreated: 1,
plansUnchanged: 1,
plansSuperseded: 1,
plansSuppressed: 1,
plansClosed: 1, // ⭐ 只关 p9;若变成 5 说明分批把其他批的患者误判成"信号消失"
});
if (baseline === null) baseline = snap;
else expect(snap).toEqual(baseline);
});
test('⛔ 0 = single-shot 回退开关(等价于改造前的不分批行为)', () => {
expect(resolvePlanBatchSize({ PAC_PLAN_BATCH_SIZE: '0' } as NodeJS.ProcessEnv)).toBe(0);
});
test('⛔ 批大小夹在 [1, 20000] —— 上限守 PG 32767 bind(latest/persona 走 patientId in ids)', () => {
expect(resolvePlanBatchSize({ PAC_PLAN_BATCH_SIZE: '999999' } as NodeJS.ProcessEnv)).toBe(20_000);
expect(resolvePlanBatchSize({ PAC_PLAN_BATCH_SIZE: '-5' } as NodeJS.ProcessEnv)).toBe(1);
expect(resolvePlanBatchSize({} as NodeJS.ProcessEnv)).toBe(2000);
expect(resolvePlanBatchSize({ PAC_PLAN_BATCH_SIZE: 'abc' } as NodeJS.ProcessEnv)).toBe(2000);
});
});
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment