Commit 889dd6f0 by luoqi

fix(persona): 画像水位对齐到 fact 层 + 写入侧节流

## 事故

画像算的是 patient_facts,stale 判定却只看 patient_transactions 的 event_seq。
reparse 重写事实但不产新事务,于是 2026-07-21「amount_cents 应收→实付」reparse
之后紧跟的全量重算被水位闸整批 noop 掉 —— 生产日志:

    manual:cli  noop=325,979  success=79,010

后果(生产 3000 抽样):reparse 前算的画像 96% 累计消费与实时事实不符,中位高估 91%。
页面上同屏两个数打架 —— 关键事实 ¥154,350(实时)vs 画像 ¥195,750(旧口径)。
而且不只是显示错:monetaryCents 要跟实时算出的 M 分位阈值比大小 → mScore 系统性
偏高 → RFM 分群整体偏「重要」→ 直接抬高召回优先级排序。

## 修法

水位对齐:新增 personas.fact_watermark = max(patient_facts.updated_at),
与 event_watermark 取「任一落后即重算」。存量行为 NULL → 判 stale,自然回填。
stale-scan 同步加这一维。

闸放宽必然带来「算了但没变」,所以补一道**写入侧的闸**(persona-diff.ts):
  unchanged 一字未变     → 不写 feature 行,只推水位
  refreshed 只有时间派生值变 → 就地刷新,不升版本(记 refreshedAt)
  semantic  真变了       → 升版本
判据比 data(剔易变键)+ evidence.factIds,不解析 description —— 同 reason-refresh 思路,
canonical 序列化抽成 common/canonical-json.ts 两边共用(jsonb 不保留键序)。

顺带:闸里查过的水位不再在正算路径重查(原来 transaction / persona 各查两遍)。

## 本地验证(friday host,13,268 患者)

  首轮回填  noop=0(旧代码这里会大批 noop)· success=48 · refreshed=11,155 · unchanged=2,065
            → 写侧节流把 13,268 次升版本压到 48 次
  再次运行  noop=13,268 全零其余,142s → 29s,幂等
  事故复现  只改 fact 不动 transaction → 触发重算,¥24,994 订正为 ¥15,996,v1 留痕

## 部署注意

patient_facts 370 万行,新索引请先手动 CREATE INDEX CONCURRENTLY 预建(迁移里是
IF NOT EXISTS,预建后即 no-op);上线后跑一次全量 recompute-persona 补 fact_watermark。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
parent 13737a51
-- 画像水位对齐到 fact 层 + 就地刷新时刻。
--
-- 背景:画像读 patient_facts,但 stale 判定只看 event_watermark(patient_transactions)。
-- reparse 重写事实、不产新事务 → 2026-07-21「amount_cents 应收→实付」reparse 之后的全量重算
-- 被水位闸 noop 掉 325,979 人(生产日志),74% 存量画像的累计消费永久停在旧口径。
ALTER TABLE "personas"
ADD COLUMN "fact_watermark" TIMESTAMPTZ(3),
ADD COLUMN "refreshed_at" TIMESTAMPTZ(3);
-- fact 水位查询(每次重算前 MAX(updated_at) per patient)的支撑索引:倒序扫首行即出,index-only。
-- 没它就得沿 (patient_id, status) 扫完该患者全部 fact 再回堆取 updated_at —— 闸本身变成开销。
-- IF NOT EXISTS:prod 先手动 CREATE INDEX CONCURRENTLY 预建(patient_facts 370 万行,免锁写),
-- 再 deploy 时本句即 no-op。预建命令:
-- CREATE INDEX CONCURRENTLY IF NOT EXISTS "patient_facts_patient_id_updated_at_idx"
-- ON "patient_facts"("patient_id", "updated_at");
CREATE INDEX IF NOT EXISTS "patient_facts_patient_id_updated_at_idx"
ON "patient_facts"("patient_id", "updated_at");
-- ⚠️ 部署后需一次性全量重算把存量 fact_watermark 补上(NULL 一律判 stale):
-- pnpm recompute-persona -- --host=jvs-dw --concurrency=8
-- 不补也能自愈 —— 每晚 stale-scan 会按 LIMIT 速率回填,只是慢。
...@@ -659,6 +659,10 @@ model PatientFact { ...@@ -659,6 +659,10 @@ model PatientFact {
/// 原仅 (patient_id, status) type 走内存 Filter,facts/患者 多时 heap fetch 超线性膨胀(13万患者全量重算瓶颈) /// 原仅 (patient_id, status) type 走内存 Filter,facts/患者 多时 heap fetch 超线性膨胀(13万患者全量重算瓶颈)
/// 实测本地 selectHits 7.3s2.0s(~3.6×),prod 收益更大。 /// 实测本地 selectHits 7.3s2.0s(~3.6×),prod 收益更大。
@@index([patientId, type, status]) @@index([patientId, type, status])
/// 画像水位:max(updated_at) per patient(PersonaService.readWatermarks 每次重算前必查)
/// 倒序扫首行即可命中 index-only,不回堆 —— 没这条要按 (patientId,status) 扫完该患者全部
/// fact 再取 heap 上的 updated_at,重患者(千级 fact)会把"闸"本身变成开销。
@@index([patientId, updatedAt])
/// 按患者拉 actual 时间轴 /// 按患者拉 actual 时间轴
@@index([patientId, occurredAt]) @@index([patientId, occurredAt])
/// 按患者查 planned 列表(待办 / 即将复查) /// 按患者查 planned 列表(待办 / 即将复查)
...@@ -695,10 +699,21 @@ model Persona { ...@@ -695,10 +699,21 @@ model Persona {
computedAt DateTime @default(now()) @map("computed_at") @db.Timestamptz(3) computedAt DateTime @default(now()) @map("computed_at") @db.Timestamptz(3)
/// 被新版本取代的时间;null 表示当前版本 /// 被新版本取代的时间;null 表示当前版本
supersededAt DateTime? @map("superseded_at") @db.Timestamptz(3) supersededAt DateTime? @map("superseded_at") @db.Timestamptz(3)
/// 本次重算时已响应到的 transaction.event_seq;判断 persona 是否 stale 用。 /// 就地刷新时刻(只有易变数值变 —— 天数/年龄/滑窗聚合 —— 时不升版本,原地改 feature )
/// Persona 实际读 patient_facts 当前 active 版本算 features; /// null = 从未就地刷新过。UI "更新于" refreshedAt ?? computedAt;
/// 水位是触发 / stale 用的,不是消费 transaction /// computedAt 保持"本版本创建时刻"不动 —— 覆盖率日报靠它反推历史快照,前移会算错昨日覆盖。
refreshedAt DateTime? @map("refreshed_at") @db.Timestamptz(3)
/// 本次重算时已响应到的 transaction.event_seq
/// ⚠️ 单靠它判 stale 是错的( factWatermark)—— 保留为辅助水位。
eventWatermark BigInt? @map("event_watermark") eventWatermark BigInt? @map("event_watermark")
/// 本次重算所见 max(patient_facts.updated_at) —— **画像真正的输入水位**
///
/// 画像读的是 patient_facts,但历史上只用 eventWatermark(transaction) stale
/// reparse 只重写 fact、不产新 transaction,于是 2026-07-21 那次
/// amount_cents 应收→实付」reparse 之后的全量重算被水位闸 noop 325,979 (生产实测),
/// 74% 存量画像的累计消费永久停在旧口径(中位高估 91%),还会污染 RFM M 分位打分。
/// null = 存量老版本(没记过 fact 水位) 判定视为 stale,下次必重算一遍补上。
factWatermark DateTime? @map("fact_watermark") @db.Timestamptz(3)
host Host @relation(fields: [hostId], references: [id]) host Host @relation(fields: [hostId], references: [id])
/// Restrict:画像版本流是审计 + 召回合规依据,不随 patient 物删消失 /// Restrict:画像版本流是审计 + 召回合规依据,不随 patient 物删消失
......
...@@ -114,8 +114,10 @@ async function bootstrap() { ...@@ -114,8 +114,10 @@ async function bootstrap() {
`(concurrency=${args.concurrency}${args.shard ? `, shard=${args.shard.i}/${args.shard.k}` : ''}) ...`, `(concurrency=${args.concurrency}${args.shard ? `, shard=${args.shard.i}/${args.shard.k}` : ''}) ...`,
); );
// noop = 水位闸命中(自上次重算无新事件,沿用旧 persona),是正常结果,不并入 failed。 // 五档结果(见 PersonaService.recompute):noop=闸拦下没算 / unchanged=算了内容没变 /
let ok = 0, partial = 0, noop = 0, failed = 0, done = 0; // refreshed=只有天数等易变值变(就地刷新不升版本) / ok+partial=升了版本。都不是错误。
// 看 unchanged 占比就知道闸放宽后有多少是白算的 —— 调水位口径的依据。
let ok = 0, partial = 0, noop = 0, unchanged = 0, refreshed = 0, failed = 0, done = 0;
const verbose = total <= 500; // 小批逐条日志;大批只打进度,免刷屏 const verbose = total <= 500; // 小批逐条日志;大批只打进度,免刷屏
const runOne = async (p: { id: string; externalId: string; name: string | null }) => { const runOne = async (p: { id: string; externalId: string; name: string | null }) => {
try { try {
...@@ -128,6 +130,8 @@ async function bootstrap() { ...@@ -128,6 +130,8 @@ async function bootstrap() {
if (r.status === 'success') ok++; if (r.status === 'success') ok++;
else if (r.status === 'partial') partial++; else if (r.status === 'partial') partial++;
else if (r.status === 'noop') noop++; else if (r.status === 'noop') noop++;
else if (r.status === 'unchanged') unchanged++;
else if (r.status === 'refreshed') refreshed++;
else failed++; else failed++;
} catch (err) { } catch (err) {
failed++; failed++;
...@@ -143,7 +147,10 @@ async function bootstrap() { ...@@ -143,7 +147,10 @@ async function bootstrap() {
} }
logger.log('─────────────────────────────────'); logger.log('─────────────────────────────────');
logger.log(`Done: success=${ok} partial=${partial} noop=${noop} failed=${failed}`); logger.log(
`Done: success=${ok} partial=${partial} refreshed=${refreshed} ` +
`unchanged=${unchanged} noop=${noop} failed=${failed}`,
);
} catch (err) { } catch (err) {
new Logger('recompute-persona').error(err instanceof Error ? err.stack : String(err)); new Logger('recompute-persona').error(err instanceof Error ? err.stack : String(err));
process.exitCode = 2; process.exitCode = 2;
......
/**
* 规范化 JSON 序列化 —— **递归按键名排序**后再 stringify。
*
* ⚠️ 凡是「把库里读出来的 jsonb 和代码里新算出来的对象比一比,判断变没变」的场景都必须用它:
* Postgres `jsonb` **不保留键序**(按键长度 + 字节序重排存储),而代码构造对象的顺序是源码顺序。
* 直接 `JSON.stringify` 比较,两边键序必然不同 → 每次都判"变了" → 每轮重算全量重写。
* (plan_reasons 实测:不加这层,连跑三次每次都刷新。)
*
* 现有消费方:
* - plan/engine/reason-refresh.ts 召回原因就地刷新判定
* - persona/persona-diff.ts 画像写入判定
*/
export function canonicalJson(v: unknown): string {
return JSON.stringify(normalize(v));
}
/** 递归:数组保序(顺序有语义),对象按键名排序,undefined 值丢弃(对齐 JSON.stringify 语义)。 */
function normalize(x: unknown): unknown {
if (Array.isArray(x)) return x.map(normalize);
if (x !== null && typeof x === 'object') {
const o = x as Record<string, unknown>;
return Object.keys(o)
.sort()
.reduce<Record<string, unknown>>((acc, k) => {
if (o[k] !== undefined) acc[k] = normalize(o[k]);
return acc;
}, {});
}
return x;
}
import { canonicalJson } from '../../common/canonical-json';
import type { PersonaFeatureDraft } from './features/feature.interface';
/**
* 画像写入判定 —— 「算完之后要不要写、写成新版本还是就地刷新」。
*
* ── 为什么需要 ──
* 原逻辑是「闸没拦住 → 无条件写新版本 + 全套 feature 行」。闸(水位)因此被迫收得很紧,
* 而收紧的代价是 **2026-07-21 reparse 后全量重算被 noop 掉 325,979 人**(生产实测):
* reparse 只改 patient_facts、不动 patient_transactions,而水位盯的是 transaction。
* 结果 74% 的存量画像把「累计消费」永久停在旧口径(应收),中位高估 91%,
* 且这个数会去和实时算出的 M 分位阈值比大小 → mScore 系统性偏高 → 召回优先级被抬。
*
* 修法是把水位对齐到 fact 层(见 PersonaService.readWatermarks)。闸放宽之后,
* 「算了但内容没变」的情况会变多 —— 本模块就是那道**写入侧的闸**:内容没变就不写,
* 只有易变数值变了就地刷新,真变了才升版本。省下的是版本行 + 3.1 GB persona_features 的增长。
*
* ── 判据:比 data(去掉易变键)+ evidence.factIds,**不比 description** ──
* description 完全由 data 派生(`${zh} · 末诊${lastDays}天 · 就诊${totalVisits}次`),
* 所以比 data 就够;不解析文案 —— 将来改文案模板不会引发误判。同 reason-refresh 的思路。
*/
export type PersonaWriteVerdict =
/** 内容一字未变 → 不写 feature 行,只推水位 */
| 'unchanged'
/** 只有易变数值变(纯时间流逝导致)→ 就地刷新 description/data,不升版本 */
| 'volatile'
/** 语义变了(标签/取值/证据集)→ 升版本,supersede 旧版 */
| 'semantic';
/** 库里读出来的 persona_features 行(比较用的最小投影) */
export interface StoredPersonaFeature {
key: string;
description: string;
score: number | null;
data: unknown;
evidence: unknown;
}
/**
* 易变键 —— **判据:拿同一批事实、换个时刻重算,值就会变**。
*
* 这类键的变化不代表患者变了,只代表"时间过去了",不该升版本。分三类:
* ① 纯时钟:距今天数 / 周岁
* ② 时间做分母:年均就诊次数
* ③ 滑动窗口聚合:近 1 年消费额、近 1 年履约/迟到率、近 2 年预约时段占比
*
* ⚠️ 加新的时间派生 data 键时必须登记到这里,否则每次重算都会白升一个版本。
* 对应回归见 tests/persona-diff.spec.ts「易变键清单」。
*
* ⚠️ 反过来也要小心:**不要**把 monetaryCents / totalVisits / types 这类"事实变才变"的键
* 塞进来 —— 那正是 E 事故要修的东西,漏判就等于修了个寂寞。
*/
export const VOLATILE_DATA_KEYS: ReadonlySet<string> = new Set([
// ① 纯时钟
'recencyDays', // rfm / lifecycle_stage:距末诊天数
'firstSeenDays', // rfm / lifecycle_stage:距首诊天数
'lastVisitDays', // urgency_level:距末诊天数
'daysSince', // potential_treatment.detail[]:诊断距今天数
'ageYears', // age_bracket:周岁(跨档会改 bracket,那才是语义变化)
// ② 时间做分母
'visitsPerYear', // lifecycle_stage:年均就诊
// ③ 滑动窗口聚合
'netThisCents', // lifecycle_stage:近 1 年净消费
'fulfillRate', // special_attention:近 1 年履约率
'apptDecided', // special_attention:近 1 年已决预约数
'lateRate', // special_attention:近 1 年迟到率
'arrivalRecs', // special_attention:近 1 年到店记录数
'recordCount', // time_preference:近 2 年预约数
'weekdayPct', // time_preference:近 2 年工作日占比
'weekendPct',
'morningPct',
'afternoonPct',
'eveningPct',
]);
/** 递归剔除易变键(嵌套对象 / 数组元素同样处理,如 potential_treatment.detail[].daysSince) */
function stripVolatile(v: unknown): unknown {
if (Array.isArray(v)) return v.map(stripVolatile);
if (v !== null && typeof v === 'object') {
const o = v as Record<string, unknown>;
const out: Record<string, unknown> = {};
for (const k of Object.keys(o)) {
if (VOLATILE_DATA_KEYS.has(k)) continue;
out[k] = stripVolatile(o[k]);
}
return out;
}
return v;
}
/** evidence.factIds 排序后取出;批次间顺序会抖,不排序会误判成变化 */
function factIdsOf(evidence: unknown): string[] {
if (evidence == null || typeof evidence !== 'object') return [];
const ids = (evidence as { factIds?: unknown }).factIds;
return Array.isArray(ids) ? [...ids].map(String).sort() : [];
}
/** 语义指纹:key + score + data(去易变)+ 证据集。description 不参与(派生物)。 */
function semanticKey(f: StoredPersonaFeature | PersonaFeatureDraft): string {
return canonicalJson({
key: f.key,
score: f.score ?? null,
data: stripVolatile(f.data ?? null),
factIds: factIdsOf(f.evidence),
});
}
/** 全量指纹:语义指纹 + 易变值 + 文案。用来区分"完全没变"和"只有天数变了"。 */
function fullKey(f: StoredPersonaFeature | PersonaFeatureDraft): string {
return canonicalJson({
key: f.key,
description: f.description,
score: f.score ?? null,
data: f.data ?? null,
factIds: factIdsOf(f.evidence),
});
}
/** 按 key 排序后逐条拼 —— 特征集顺序在库里/代码里不一致,不排序会误判 */
function fingerprint(
features: ReadonlyArray<StoredPersonaFeature | PersonaFeatureDraft>,
of: (f: StoredPersonaFeature | PersonaFeatureDraft) => string,
): string {
return features
.map(of)
.sort()
.join('\n');
}
/**
* 该患者这次算出来的画像,相对当前 active 版本要怎么落库。
*
* @param current 当前 active persona 的 feature 行;null / 空数组 = 没有当前版本 → 必写(semantic)
* @param next 本次算出的 drafts
*/
export function personaWriteVerdict(
current: readonly StoredPersonaFeature[] | null,
next: readonly PersonaFeatureDraft[],
): PersonaWriteVerdict {
// 首次(无 active 版本)必写。注意:current 是空数组但 persona 行存在的情况也走这里 ——
// 只有"上一版一个特征都没产出"才会这样,极罕见,多写一版无害。
if (!current || current.length === 0) return 'semantic';
if (next.length === 0) return 'semantic'; // 从有到无(全部降级)也是语义变化
if (fingerprint(current, semanticKey) !== fingerprint(next, semanticKey)) return 'semantic';
if (fingerprint(current, fullKey) !== fingerprint(next, fullKey)) return 'volatile';
return 'unchanged';
}
...@@ -19,30 +19,7 @@ ...@@ -19,30 +19,7 @@
* 据此判定既充分,又不依赖文案格式(将来 scenario 改文案模板不会引发误刷)。 * 据此判定既充分,又不依赖文案格式(将来 scenario 改文案模板不会引发误刷)。
*/ */
/** import { canonicalJson as canonical } from '../../../common/canonical-json';
* 规范化序列化:**递归按键名排序**后再 stringify。
*
* ⚠️ 必须这么做 —— Postgres `jsonb` **不保留键序**(按键长度+字节序重排存储),
* 而代码里构造 signals 的顺序是源码顺序。直接 JSON.stringify 比较,库里读出来的
* 和新算出来的键序必然不同 → 每次都判"变了" → 每日重算把全部 plan_reasons 重写一遍,
* 恰好是本模块要避免的事(本地实测:不加这层,连跑三次每次都刷新)。
*/
function canonical(v: unknown): string {
const norm = (x: unknown): unknown => {
if (Array.isArray(x)) return x.map(norm);
if (x !== null && typeof x === 'object') {
const o = x as Record<string, unknown>;
return Object.keys(o)
.sort()
.reduce<Record<string, unknown>>((acc, k) => {
if (o[k] !== undefined) acc[k] = norm(o[k]);
return acc;
}, {});
}
return x;
};
return JSON.stringify(norm(v));
}
/** 稳定投影:剔除随时间必变、不代表语义变化的字段 */ /** 稳定投影:剔除随时间必变、不代表语义变化的字段 */
function stableSignals(signals: unknown): unknown { function stableSignals(signals: unknown): unknown {
......
...@@ -43,9 +43,12 @@ export class PersonaRecomputeProcessor extends WorkerHost { ...@@ -43,9 +43,12 @@ export class PersonaRecomputeProcessor extends WorkerHost {
`persona=${r.personaId ?? 'noop'} v${r.version} features=${r.featureCount} duration=${r.durationMs}ms`, `persona=${r.personaId ?? 'noop'} v${r.version} features=${r.featureCount} duration=${r.durationMs}ms`,
); );
// 链式触发:只在画像确有更新时(success / partial)才下推 plan 重算 // 链式触发:只在画像**语义**确有更新时(success / partial)才下推 plan 重算。
// noop = 水位幂等(没新事件)→ plan 也不会变,省一次 // noop 水位幂等,压根没算 → plan 也不会变
// failed = persona 写库挂了 → 不传播,等重试 / 死信 // unchanged 算了但内容一字未变 → 同上
// refreshed 只有天数/年龄这类易变值变了 → 召回打分读的是 rfm.segment / urgency.level /
// lifecycle.stage 这些**分档**,分档没变就不必重排(变了会判成 semantic)
// failed persona 写库挂了 → 不传播,等重试 / 死信
let planEnqueued = false; let planEnqueued = false;
if (r.status === 'success' || r.status === 'partial') { if (r.status === 'success' || r.status === 'partial') {
try { try {
...@@ -64,7 +67,9 @@ export class PersonaRecomputeProcessor extends WorkerHost { ...@@ -64,7 +67,9 @@ export class PersonaRecomputeProcessor extends WorkerHost {
} }
return { return {
ok: r.status === 'success' || r.status === 'noop', // ok = 干净完成。新增的 unchanged / refreshed 与 noop 同档(都属"无需升版本"的正常结果);
// partial 仍不算 ok(有 extractor 抛错),沿用原语义。
ok: r.status === 'success' || r.status === 'noop' || r.status === 'unchanged' || r.status === 'refreshed',
personaId: r.personaId ?? null, personaId: r.personaId ?? null,
status: r.status, status: r.status,
planEnqueued, planEnqueued,
......
...@@ -11,8 +11,15 @@ import { QueueProducer } from './queue-producer.service'; ...@@ -11,8 +11,15 @@ import { QueueProducer } from './queue-producer.service';
* - 应用错误期间漏入队的(部署 / panic 错过 enqueue) * - 应用错误期间漏入队的(部署 / panic 错过 enqueue)
* - cold-import 进来但没触发 enqueue 的(脚本路径) * - cold-import 进来但没触发 enqueue 的(脚本路径)
* *
* 判定:max(patient_transactions.event_seq) > personas.event_watermark (或 persona 不存在) * 判定(两个水位任一落后即算 stale,或 persona 根本不存在):
* → 这个 patient 的画像落后,enqueue persona-recompute(BullMQ jobId 自动去重) * - max(patient_transactions.event_seq) > personas.event_watermark
* - max(patient_facts.updated_at) > personas.fact_watermark ← 主水位
* → enqueue persona-recompute(BullMQ jobId 自动去重)
*
* ⚠️ fact 水位是 2026-07 补的:画像读 patient_facts,但原来只判 transaction,
* reparse(重写事实、不产事务)后的重算被全部 noop 掉。存量 persona 的 fact_watermark
* 为 NULL → 本扫描会持续判它们 stale,以每晚 LIMIT 的速率**自然回填**。
* 要一次性补完请跑 `pnpm recompute-persona -- --host=<h> --concurrency=N`。
* *
* 安全:扫描走 GROUP BY + HAVING,不锁表。同一 patient 即便已在主路径队列里, * 安全:扫描走 GROUP BY + HAVING,不锁表。同一 patient 即便已在主路径队列里,
* jobId 去重保证不会重复算。 * jobId 去重保证不会重复算。
...@@ -74,28 +81,39 @@ export class StaleScanService { ...@@ -74,28 +81,39 @@ export class StaleScanService {
/** /**
* 找出"画像 stale"的患者: * 找出"画像 stale"的患者:
* - 有 transaction 但没 persona(从来没算过) * - 有 transaction 但没 persona(从来没算过)
* - 有 persona 但 watermark < latest event_seq(漏了更新) * - 有 persona 但 event 水位落后(漏了更新)
* - 有 persona 但 **fact 水位落后**(reparse / 事实层订正 —— 不产新事务的那类变更)
* *
* 用一句 raw SQL 兼顾种情况;返回 (hostId, tenantId, patientId) 三元组。 * 用一句 raw SQL 兼顾种情况;返回 (hostId, tenantId, patientId) 三元组。
*/ */
private async findStalePatients(): Promise< private async findStalePatients(): Promise<
Array<{ hostId: string; tenantId: string; patientId: string }> Array<{ hostId: string; tenantId: string; patientId: string }>
> { > {
// Postgres:LEFT JOIN active persona,WHERE 取 watermark IS NULL 或 max_seq > watermark
const rows = await this.prisma.$queryRaw< const rows = await this.prisma.$queryRaw<
Array<{ host_id: string; tenant_id: string; patient_id: string }> Array<{ host_id: string; tenant_id: string; patient_id: string }>
>` >`
SELECT t.host_id, t.tenant_id, t.patient_id WITH tx AS (
FROM (
SELECT host_id, tenant_id, patient_id, MAX(event_seq) AS max_seq SELECT host_id, tenant_id, patient_id, MAX(event_seq) AS max_seq
FROM patient_transactions FROM patient_transactions
WHERE patient_id IS NOT NULL WHERE patient_id IS NOT NULL
GROUP BY host_id, tenant_id, patient_id GROUP BY host_id, tenant_id, patient_id
) t ), fx AS (
-- (patient_id, updated_at) 索引上的 index-only 聚合,不回堆
SELECT patient_id, MAX(updated_at) AS max_fact_at
FROM patient_facts
GROUP BY patient_id
)
SELECT tx.host_id, tx.tenant_id, tx.patient_id
FROM tx
LEFT JOIN personas p LEFT JOIN personas p
ON p.patient_id = t.patient_id ON p.patient_id = tx.patient_id
AND p.superseded_at IS NULL AND p.superseded_at IS NULL
WHERE p.event_watermark IS NULL OR p.event_watermark < t.max_seq LEFT JOIN fx ON fx.patient_id = tx.patient_id
WHERE p.id IS NULL
OR p.event_watermark IS NULL
OR p.event_watermark < tx.max_seq
OR (fx.max_fact_at IS NOT NULL
AND (p.fact_watermark IS NULL OR p.fact_watermark < fx.max_fact_at))
LIMIT 5000 LIMIT 5000
`; `;
return rows.map((r) => ({ return rows.map((r) => ({
......
import {
personaWriteVerdict,
VOLATILE_DATA_KEYS,
type StoredPersonaFeature,
} from '../src/modules/persona/persona-diff';
import type { PersonaFeatureDraft } from '../src/modules/persona/features/feature.interface';
/**
* 画像写入判定回归。
*
* 背景(生产事故):画像读 patient_facts,stale 判定却只看 patient_transactions 的水位。
* 2026-07-21「amount_cents 应收→实付」reparse 只改事实、不产事务 → 之后的全量重算
* 被水位闸 noop 掉 325,979 人,74% 存量画像的累计消费永久停在旧口径(中位高估 91%),
* 而这个数会去和实时算的 M 分位阈值比大小 → mScore 系统性偏高 → 召回优先级被抬。
*
* 修法是把水位对齐到 fact 层。闸放宽之后「算了但没变」变多,本模块是**写入侧的闸**:
* unchanged 一字未变 → 不写 feature 行
* volatile 只有时间派生值变 → 就地刷新,不升版本
* semantic 真变了 → 升版本
*
* 本文件同时守住两条红线:
* ① 不能把"事实变才变"的键(monetaryCents 等)误判成易变 —— 那就等于没修事故
* ② 不能把纯时间漂移judg成语义变化 —— 那会让每次重算都白升一版
*/
const draft = (over: Partial<PersonaFeatureDraft> = {}): PersonaFeatureDraft =>
({
key: 'rfm',
description: '重要价值 · 距上次 140 天 · 就诊20次 · 累计¥154,350',
score: null,
data: {
segment: 'important_value',
rScore: 5,
fScore: 5,
mScore: 5,
recencyDays: 140,
firstSeenDays: 3000,
freqCount: 20,
monetaryCents: 15435000,
valueTier: 4,
riskScore: 2,
hasTreatmentGap: true,
},
evidence: { factIds: ['f1', 'f2'] },
...over,
}) as PersonaFeatureDraft;
/** draft → 库里那行的形状(字段同名,类型放宽) */
const stored = (d: PersonaFeatureDraft): StoredPersonaFeature => ({
key: d.key,
description: d.description,
score: d.score ?? null,
data: d.data ?? null,
evidence: d.evidence,
});
const verdict = (cur: PersonaFeatureDraft[], next: PersonaFeatureDraft[]) =>
personaWriteVerdict(cur.map(stored), next);
describe('personaWriteVerdict — 三档判定', () => {
test('⭐ 完全相同 → unchanged(不写任何 feature 行)', () => {
expect(verdict([draft()], [draft()])).toBe('unchanged');
});
test('⭐ 只有天数漂移(recencyDays 140→151)→ volatile,**不升版本**', () => {
const next = draft({
description: '重要价值 · 距上次 151 天 · 就诊20次 · 累计¥154,350',
data: { ...(draft().data as object), recencyDays: 151 },
});
expect(verdict([draft()], [next])).toBe('volatile');
});
test('⭐⭐ 累计消费变了(应收→实付,正是事故本体)→ semantic,必须升版本', () => {
const next = draft({
description: '重要价值 · 距上次 140 天 · 就诊20次 · 累计¥195,750',
data: { ...(draft().data as object), monetaryCents: 19575000 },
});
expect(verdict([draft()], [next])).toBe('semantic');
});
test('⭐ 分群变了(mScore 掉档导致 segment 变)→ semantic', () => {
const next = draft({
data: { ...(draft().data as object), segment: 'general_value', mScore: 3 },
});
expect(verdict([draft()], [next])).toBe('semantic');
});
test('首次(无当前版本)→ semantic', () => {
expect(personaWriteVerdict(null, [draft()])).toBe('semantic');
expect(personaWriteVerdict([], [draft()])).toBe('semantic');
});
test('特征全没了(从有到无)→ semantic', () => {
expect(verdict([draft()], [])).toBe('semantic');
});
test('新增一个特征 → semantic', () => {
const extra = draft({ key: 'age_bracket', data: { bracket: 'senior', label: '老年', ageYears: 58 } });
expect(verdict([draft()], [draft(), extra])).toBe('semantic');
});
test('掉了一个特征(降级为不打标签)→ semantic', () => {
const extra = draft({ key: 'age_bracket', data: { bracket: 'senior', label: '老年', ageYears: 58 } });
expect(verdict([draft(), extra], [draft()])).toBe('semantic');
});
test('特征顺序不同 → 不算变化(库里/代码里顺序本就不保证)', () => {
const a = draft({ key: 'age_bracket', data: { bracket: 'senior' } });
expect(verdict([draft(), a], [a, draft()])).toBe('unchanged');
});
});
describe('易变键清单 —— 时间派生值不该升版本', () => {
// 判据:拿同一批事实、换个时刻重算,值就会变
test.each([
['recencyDays', 'rfm / lifecycle_stage 距末诊天数'],
['firstSeenDays', 'rfm / lifecycle_stage 距首诊天数'],
['lastVisitDays', 'urgency_level 距末诊天数'],
['ageYears', 'age_bracket 周岁'],
['visitsPerYear', 'lifecycle_stage 年均就诊(时间做分母)'],
['netThisCents', 'lifecycle_stage 近 1 年净消费(滑窗)'],
['fulfillRate', 'special_attention 近 1 年履约率(滑窗)'],
['recordCount', 'time_preference 近 2 年预约数(滑窗)'],
])('%s 变化 → volatile(%s)', (key) => {
const base = draft({ data: { segment: 'important_value', [key]: 1 } });
const next = draft({ data: { segment: 'important_value', [key]: 999 }, description: '改了' });
expect(verdict([base], [next])).toBe('volatile');
});
test('⭐ 嵌套里的 daysSince 也要剔(potential_treatment.detail[])', () => {
const mk = (d: number) =>
draft({
key: 'potential_treatment',
description: '潜在种植 / 潜在修复',
data: {
types: ['implant', 'restoration'],
labels: ['潜在种植', '潜在修复'],
detail: [{ key: 'implant', teeth: ['16'], daysSince: d, confidence: 1 }],
},
});
expect(verdict([mk(140)], [mk(151)])).toBe('volatile');
});
test('⭐ 红线:牙位变了(detail[].teeth)→ semantic,不能被 detail 的剔除逻辑连坐', () => {
const mk = (teeth: string[]) =>
draft({
key: 'potential_treatment',
data: { types: ['implant'], detail: [{ key: 'implant', teeth, daysSince: 140 }] },
});
expect(verdict([mk(['16'])], [mk(['16', '17'])])).toBe('semantic');
});
test('⭐ 红线:清单里不能混进"事实变才变"的键', () => {
for (const k of ['monetaryCents', 'netTotalCents', 'totalVisits', 'freqCount', 'types', 'segment', 'stage', 'level', 'bracket']) {
expect(VOLATILE_DATA_KEYS.has(k)).toBe(false);
}
});
});
describe('evidence 维度', () => {
test('证据集变了(命中新 fact)→ semantic', () => {
const next = draft({ evidence: { factIds: ['f1', 'f2', 'f3'] } });
expect(verdict([draft()], [next])).toBe('semantic');
});
test('⭐ factIds 只是顺序不同 → 不算变化(批次间顺序会抖)', () => {
const next = draft({ evidence: { factIds: ['f2', 'f1'] } });
expect(verdict([draft()], [next])).toBe('unchanged');
});
test('evidence 形状异常不炸', () => {
const a = draft({ evidence: null as never });
const b = draft({ evidence: {} as never });
expect(verdict([a], [b])).toBe('unchanged');
});
});
describe('⭐ 红线:jsonb 键序不能被当成变化', () => {
// Postgres jsonb 不保留键序 —— 不规范化的话每次重算都判"变了",全量重写(3.1 GB 表)
test('顶层键序不同 → unchanged', () => {
const asStored = stored(draft());
asStored.data = {
hasTreatmentGap: true, monetaryCents: 15435000, valueTier: 4, segment: 'important_value',
riskScore: 2, freqCount: 20, firstSeenDays: 3000, recencyDays: 140,
mScore: 5, fScore: 5, rScore: 5,
};
expect(personaWriteVerdict([asStored], [draft()])).toBe('unchanged');
});
test('嵌套对象键序不同 → unchanged', () => {
const mk = (o: Record<string, unknown>) =>
draft({ key: 'potential_treatment', data: { detail: [o] } });
const a = mk({ key: 'implant', teeth: ['16'], confidence: 1 });
const b = mk({ confidence: 1, teeth: ['16'], key: 'implant' });
expect(verdict([a], [b])).toBe('unchanged');
});
test('数组顺序仍算变化(types 首项 = 主推,有序)', () => {
const a = draft({ key: 'treatment_history', data: { types: ['implant_history', 'perio_history'] } });
const b = draft({ key: 'treatment_history', data: { types: ['perio_history', 'implant_history'] } });
expect(verdict([a], [b])).toBe('semantic');
});
});
describe('红线:判定不依赖 description 文案', () => {
test('⭐ data 相同、只有文案模板改了 → volatile(就地刷新),不升版本', () => {
const next = draft({ description: '【重要价值】140 天未到诊,累计 154,350 元' });
expect(verdict([draft()], [next])).toBe('volatile');
});
});
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