Commit 4bd77113 by luoqi

merge: fix/push-schema-violation-isolation → main(脏数据不连坐 + 第三种退费表达 + 回执可自诊断)

parents f006448f 9c4065d8
Pipeline #3472 failed in 0 seconds
......@@ -234,10 +234,33 @@ const sig = require('crypto').createHmac('sha256', SECRET).update(`${ts}.${body
| 字段 | 含义 | 大于 0 时 |
|---|---|---|
| `failed` | 本批失败行数 | 零星可忽略,持续升高联系 PAC |
| `failed` | 本批**没进 PAC** 的行数 | 零星可忽略,持续升高联系 PAC |
| `factsFailed` | 行进了、但**派生不出可用事实**的条数 | 看 `factRejects` 定位到自己源库那一行 |
| `mappingMisses` | 存在未映射的原值 | 暂停推送,发样本给 PAC |
| `suspectFields` | 与历史相比列发生增减 | 确认表结构是否变更,同步 PAC |
`factsFailed > 0` 时,`factRejects` 给出定位明细(最多 10 条,同一种错通常整批同因):
```json
{
"transactionsWritten": 199,
"factsFailed": 1,
"factRejects": [
{
"type": "payment_record",
"subjectId": "payment_record:be67598c35ec11eebdaf00163e342921",
"transactionId": "940e83e0-1dd8-451c-a494-c15e1512eac5",
"issues": ["amount_cents: Too small: expected number to be >=0"]
}
]
}
```
`subjectId` 的冒号后面就是**贵方的主键**,据此回源库定位。上例是一张实收为负的红冲结算单
——`failed` 为 0(行本身收下了),但它派生不出合法的消费事实。
出于隐私,`factRejects` 只给定位信息与字段级原因,**不回传记录原文**。
### 9.2 补漏
每晚将最近 N 天数据全量重推一遍:已推的跳过、漏推的补上、变更的升版本。
......
......@@ -303,10 +303,20 @@ transforms:
where:
status: { in: ['1', '3'] }
receivable_this: { gte: 0 } # 应收≥0(迁自导出侧 WHERE;数值算子)
net_receipts_this: { gte: 0 } # ⭐ 实收≥0 —— 见 S.2 第③种表达,负实收归退费不归消费
# ── S.2 退费①:整单冲减 —— PAC 侧切分(不再靠导出 WHERE)。两种表达 union:
# ① status=4 整单反向;② status=3 且 receivable_this<0(2ac96f6d 品牌用负额表达退费)。
# ── S.2 退费①:整单冲减 —— PAC 侧切分(不再靠导出 WHERE)。**三种**表达 union:
# ① status=4 整单反向;
# ② status=3 且 receivable_this<0(2ac96f6d 品牌用负应收表达退费);
# ③ status∈{1,3} 且 receivable_this≥0 但 net_receipts_this<0 —— 应收不动、只把钱退回去。
# 金额负值由 refund parser Math.abs 归一为正 cents。
#
# ⚠️ ③ 是 2026-07-29 线上事故补的,别当冗余删掉:
# 漏了它 → 这些行落进 payment_rows → payment_record 的 amount_cents 为负 → zod 拒收
# → bulkWrite 整批失败 → 降级逐行事务 → 单批 200 行从 3 秒涨到 63.5 秒
# → 超过宿主 60s 读超时 → 对方永远收不到成功回执、无限重推同一批(实测卡了 13 小时)。
# 测试服实测:57,939 条结算里这样的只有 **4 行**(status=3 / 应收 0 或 6000 / 实收 -500~-6000),
# 与 ①② 零重叠(全库无 status=4、无应收<0)。
- kind: filter
input: patient_settlement
output: _refund_status4
......@@ -318,8 +328,15 @@ transforms:
where:
status: { equals: '3' }
receivable_this: { lt: 0 } # 数值算子:负额退费表达
- kind: filter
input: patient_settlement
output: _refund_neg_receipt
where:
status: { in: ['1', '3'] }
receivable_this: { gte: 0 } # 与 ② 互斥(② 要求应收<0),union 不会重复计
net_receipts_this: { lt: 0 }
- kind: union
inputs: ['_refund_status4', '_refund_neg3']
inputs: ['_refund_status4', '_refund_neg3', '_refund_neg_receipt']
output: refund_full_rows
# ── S.3 退费明细②:行级退费(patient_settlement_spec is_refund=1)——PAC 侧 filter ──
......
......@@ -38,6 +38,14 @@ import {
} from './clickhouse-source.service';
import { AlertService } from '../../../common/alerting/alert.service';
import { TransformEngine } from '../transforms/transform-engine';
import type { FactReject } from '../pipeline/fact-writer.service';
/**
* 单次摄入回执里最多带几条 schema 拒收明细。
* 回执是给宿主工程师看的诊断线索,不是日志转储 —— 同一种错通常整批同因,
* 给几条足够定位;全量在 PAC 服务端日志里(每条都有 [schema-violation] 行)。
*/
const FACT_REJECT_SAMPLE_CAP = 10;
import type { TransformOp } from '../transforms/transforms.schema';
import {
buildPushLookupFallbackRows,
......@@ -2054,6 +2062,9 @@ export class ColdImportService {
stats.factsUnchanged += metrics.factsUnchanged;
stats.factsEvidenceAppended += metrics.factsEvidenceAppended;
stats.factsFailed += metrics.factsFailed;
for (const r of metrics.factRejects) {
if (stats.factRejects.length < FACT_REJECT_SAMPLE_CAP) stats.factRejects.push(r);
}
}
/// 兜底:createMany 整批失败时降级 per-row(慢但稳)
......@@ -2351,6 +2362,7 @@ export class ColdImportService {
factsUnchanged: 0,
factsEvidenceAppended: 0,
factsFailed: 0,
factRejects: [],
};
}
......@@ -2367,6 +2379,7 @@ export class ColdImportService {
factsUnchanged: 0,
factsEvidenceAppended: 0,
factsFailed: 0,
factRejects: [],
sampleCanonical: [],
mappingMisses: [],
suspectFields: [],
......@@ -2388,6 +2401,9 @@ interface TotalsBlock {
factsUnchanged: number;
factsEvidenceAppended: number;
factsFailed: number;
/// schema 拒收明细(定位信息,无 content 原文)—— push 回执透给宿主自诊断。
/// 累计上限见 FACT_REJECT_SAMPLE_CAP:回执是给人看的,不是日志转储。
factRejects: FactReject[];
}
/// 沿 transforms 的 output→input 链回溯,找到非任何 transform 产出的"原始源表"名。
......
......@@ -214,10 +214,17 @@ export class FactWriter {
*
* 失败处理:
* - 整批 $transaction 失败 → 全部回滚 → 抛 BulkWriteFailedError(调用方降级 per-entry)
* - zod 失败 → 同步抛(整批拒绝;调用方拆 entry 跑 writeDraft 兜底)
* - **zod 失败 → 只踢掉违规的那几条**(rejected),其余照常走 bulk。
*
* ⚠️ zod 失败**曾经**是整批抛错的,2026-07-29 线上事故后改成隔离,别改回去:
* FRIDAY 推的 57,939 条结算里有 4 行负实收(退费的第三种表达,manifest 当时没覆盖),
* payment_record 的 amount_cents≥0 校验拒收 → 整批抛 → ParserPipeline 降级 per-entry →
* **每 fact 一个 $transaction** → 单批 200 行从 3 秒涨到 63.5 秒 → 超过宿主 60s 读超时 →
* 对方收不到成功回执、无限重推同一批,实测卡了 13 小时一条没进。
* 一行脏数据不该让同批另外 199 行陪葬 —— 隔离它、如实回报它,批照常走。
*/
async bulkWrite(entries: BulkEntry[]): Promise<FactWriteResult[]> {
if (entries.length === 0) return [];
async bulkWrite(entries: BulkEntry[]): Promise<BulkWriteOutcome> {
if (entries.length === 0) return { results: [], rejected: [] };
// 假设全 batch 同 hostId+tenantId(processSubject 是按 host+tenant 跑的,符合)
// 防御:校验所有 entry 同 host+tenant
......@@ -230,11 +237,32 @@ export class FactWriter {
}
}
// 1. 全部 zod 校验
const validatedEntries = entries.map((e) => {
const content = validateFactContent(e.draft.type, e.draft.subjectId, e.draft.content) as Prisma.InputJsonValue;
return { ...e, validatedContent: content, hash: this.hashContent(content) };
// 1. 逐条 zod 校验 —— 违规的踢进 rejected,不连坐同批其余行(见上方 ⚠️)
const validatedEntries: Array<BulkEntry & { validatedContent: Prisma.InputJsonValue; hash: string }> = [];
const rejected: FactReject[] = [];
for (const e of entries) {
try {
const content = validateFactContent(
e.draft.type,
e.draft.subjectId,
e.draft.content,
) as Prisma.InputJsonValue;
validatedEntries.push({ ...e, validatedContent: content, hash: this.hashContent(content) });
} catch (err) {
if (!(err instanceof FactContentSchemaError)) throw err; // 非 schema 错(编程错)不吞
rejected.push({
type: e.draft.type,
subjectId: e.draft.subjectId,
transactionId: e.transactionId,
issues: err.issues.map((i) => `${i.path.join('.') || '(root)'}: ${i.message}`),
});
this.logger.error(
`[schema-violation] type=${e.draft.type} subjectId=${e.draft.subjectId} ` +
`transactionId=${e.transactionId} issues=${JSON.stringify(err.issues)}`,
);
}
}
if (validatedEntries.length === 0) return { results: [], rejected };
// 2. 一次 SELECT 把所有相关 subject 的 latest version 拿回
// ⭐ 集团内 subject_id 跨品牌会撞 → 史表查带 source_unit 过滤,版本 map 按 (source_unit, subject_id) 建键
......@@ -413,7 +441,7 @@ export class FactWriter {
);
}
return results;
return { results, rejected };
}
/**
......@@ -485,6 +513,28 @@ export interface FactWriteResult {
version: number;
}
/**
* 被 fact.content schema 拒收的单条 draft —— 一路透到 push 回执,让宿主能自诊断。
*
* 只带定位信息(哪张源记录、哪个字段、为什么),**不带 content 原文**:
* 回执会进对方日志,原文里有患者姓名/金额等,不该顺着回执外流。
*/
export interface FactReject {
/// fact 类型(= assembler canonical,如 payment_record)
type: string;
/// 出问题的 subject(形如 payment_record:<宿主主键>)—— 宿主据此定位到具体那一行
subjectId: string;
transactionId: string;
/// 人读的字段级原因,如 "amount_cents: Too small: expected number to be >=0"
issues: string[];
}
/// bulkWrite 结果:成功写入的 + 被 schema 拒收隔离的(不抛错,不连坐)
export interface BulkWriteOutcome {
results: FactWriteResult[];
rejected: FactReject[];
}
/// FactWriter.bulkWrite 入参条目 — 一个 draft + 它所属的 patient / host / tenant / tx 上下文
export interface BulkEntry {
draft: FactDraft;
......
......@@ -5,9 +5,11 @@ import { ParserRegistry } from './parsers/parser.registry';
import { ClinicalSignalService } from '../../clinical-signals/clinical-signal.service';
import {
type BulkEntry,
type FactReject,
FactWriter,
FactWriteResult,
} from './fact-writer.service';
import { FactContentSchemaError } from './fact-content-schemas';
/**
* ParserPipeline — transaction → fact 衍生编排器
......@@ -69,6 +71,7 @@ export class ParserPipeline {
factsEvidenceAppended: 0,
factsStaleSkipped: 0,
factsFailed: 0,
factRejects: [],
writes: [],
};
......@@ -160,6 +163,7 @@ export class ParserPipeline {
factsEvidenceAppended: 0,
factsStaleSkipped: 0,
factsFailed: 0,
factRejects: [],
writes: [],
};
......@@ -215,8 +219,16 @@ export class ParserPipeline {
if (bulkEntries.length === 0) return metrics;
// 2. 一次 bulk write,失败降级 per-entry(写一份保证收尾)
// 注:fact.content schema 违规**不再**触发降级 —— bulkWrite 内部把违规条目隔离进 rejected,
// 其余照常批量写。这里的 catch 只兜住真正的批级故障(DB 异常 / 事务超时)。
// 2026-07-29 前不是这样:一行负金额就把整批打成逐行事务,200 行 3s→63.5s,
// 直接把宿主推送卡到超时重推死循环(详见 FactWriter.bulkWrite 注释)。
try {
const results = await this.writer.bulkWrite(bulkEntries);
const { results, rejected } = await this.writer.bulkWrite(bulkEntries);
for (const r of rejected) {
metrics.factsFailed++;
metrics.factRejects.push(r);
}
// 逐个 push,不用 push(...results):单资源单批已达 5 万+ 行,spread 实参压栈会炸
// RangeError(V8 上限 ~6.5万;同型 bug 曾炸 plans 的 selectHits,见 scenario 注释)。
for (const r of results) {
......@@ -275,6 +287,14 @@ export class ParserPipeline {
}
} catch (subErr) {
metrics.factsFailed++;
if (subErr instanceof FactContentSchemaError) {
metrics.factRejects.push({
type: subErr.factType,
subjectId: subErr.subjectId,
transactionId: e.transactionId,
issues: subErr.issues.map((i) => `${i.path.join('.') || '(root)'}: ${i.message}`),
});
}
this.logger.error(
`fallback writeDraft failed: tx=${e.transactionId} subject=${e.draft.subjectId} ` +
`err=${subErr instanceof Error ? subErr.message : String(subErr)}`,
......@@ -312,6 +332,9 @@ export interface PipelineRunMetrics {
/// 乱序防护跳过数(更旧 sourceUpdatedAt 晚到,未旧覆新)
factsStaleSkipped: number;
factsFailed: number;
/// 被 fact.content schema 拒收的条目明细(定位信息,无 content 原文)——
/// 一路透到 push 回执,让宿主知道是自己哪一行有问题,而不是只看到超时。
factRejects: FactReject[];
writes: FactWriteResult[];
}
......
......@@ -68,10 +68,22 @@ export class PushReceiverService {
txn: a.txn + s.transactionsWritten,
dup: a.dup + s.duplicates,
failed: a.failed + s.failed,
factsFailed: a.factsFailed + s.factsFailed,
}),
{ txn: 0, dup: 0, failed: 0 },
{ txn: 0, dup: 0, failed: 0, factsFailed: 0 },
);
// fact.content 被 schema 拒收的明细 —— 必须透出去,否则宿主只能看到"成功但没进数据"。
// 2026-07-29 事故:4 行负金额结算被拒,回执却是 success/failed=0,对方无从判断,
// 只看到超时 → 重推同一批 13 小时。定位信息(哪条 subject、哪个字段)不含 content 原文。
const factRejects = r.perResource.flatMap((s) => s.factRejects);
if (factRejects.length > 0) {
this.logger.warn(
`push/rows host=${hostName} source=${source}${agg.factsFailed} 条 fact 被 schema 拒收,` +
`样本=${JSON.stringify(factRejects.slice(0, 3))}`,
);
}
// ── 关键字段守卫:整批全失败(非重复、无一落库)→ 拒收报错,别让宿主静默丢数据 ──
// 典型成因:宿主改了身份/时间列名 → 逐行合成失败。失败行不落任何库(raw 只活在 txn),
// 拒绝是无损的;宿主修好后整批重推即可(幂等)。SyncLog 已标 FAILED + 告警(ingestRawTables 内)。
......@@ -120,6 +132,10 @@ export class PushReceiverService {
transactionsWritten: agg.txn,
duplicates: agg.dup,
failed: agg.failed,
/// transaction 已落账、但衍生 fact 被 content schema 拒收的条数(0 = 全部正常)
factsFailed: agg.factsFailed,
/// 上面那些拒收的定位明细(封顶若干条),宿主据此定位到自己源库的具体行
factRejects,
personaEnqueued,
mappingMisses: r.mappingMisses.length,
suspectFields: r.suspectFields.length,
......
......@@ -86,6 +86,26 @@ export const PushRowsResponseSchema = z.object({
transactionsWritten: z.number().int(),
duplicates: z.number().int(),
failed: z.number().int(),
/**
* transaction 已落账、但衍生 fact 被 content schema 拒收的条数。
*
* 跟 failed 分开:failed 是"这行没进 PAC";factsFailed 是"行进了,但派生不出可用事实"
* (典型:金额为负、必填字段空)。**>0 表示这批数据有一部分对 PAC 无效,要看 factRejects**。
*
* 2026-07-29 补:此前这类拒收既不计数也不回报,回执是 success/failed=0,
* 宿主完全看不出问题在自己哪一行,只看到超时 → 重推同一批 13 小时(见 factRejects)。
*/
factsFailed: z.number().int().describe('衍生 fact 被 schema 拒收的条数(0 = 全部正常)'),
factRejects: z
.array(
z.object({
type: z.string().describe('fact 类型,如 payment_record'),
subjectId: z.string().describe('形如 <类型>:<宿主主键> —— 据此定位到源库那一行'),
transactionId: z.string(),
issues: z.array(z.string()).describe('字段级原因,如 "amount_cents: 需 >= 0"'),
}),
)
.describe('拒收明细样本(封顶 10 条;不含 content 原文,避免患者信息随回执外流)'),
personaEnqueued: z.number().int().describe('入队画像重算的患者数(plan 不在此触发,由定时任务保)'),
/// 宿主自检信号:样本批推完看这两个数,>0 先停下联系 PAC。
mappingMisses: z.number().int().describe('映射覆盖缺口种类数(有原值落 _default)'),
......
......@@ -318,7 +318,7 @@ describe('FactWriter.bulkWrite | 时间锚变更检测', () => {
]);
const w = writer(prisma);
const results = await w.bulkWrite([
const { results } = await w.bulkWrite([
{ ...baseInput, draft: encounterDraft({ subjectId: 'enc-1', content: { encounter_external_id: 'enc-1' } }), transactionId: 'tx-9' },
{ ...baseInput, draft: encounterDraft({ subjectId: 'enc-2', content: { encounter_external_id: 'enc-1' }, occurredAt: T0_MINUS_8H }), transactionId: 'tx-9' },
]);
......@@ -340,7 +340,7 @@ describe('FactWriter.bulkWrite | 时间锚变更检测', () => {
const { prisma, rows } = makeStore([seedActive()]);
const w = writer(prisma);
const results = await w.bulkWrite([
const { results } = await w.bulkWrite([
{ ...baseInput, draft: encounterDraft({ occurredAt: T0_MINUS_8H }), transactionId: 'tx-2' },
{ ...baseInput, draft: encounterDraft({ occurredAt: new Date('2026-07-01T03:00:00Z') }), transactionId: 'tx-3' },
]);
......@@ -354,7 +354,7 @@ describe('FactWriter.bulkWrite | 时间锚变更检测', () => {
const { prisma, rows } = makeStore([
seedActive({ sourceUpdatedAt: new Date('2026-07-02T00:00:00Z') }),
]);
const results = await writer(prisma).bulkWrite([
const { results } = await writer(prisma).bulkWrite([
{
...baseInput,
draft: encounterDraft({ occurredAt: T0_MINUS_8H }),
......@@ -366,3 +366,67 @@ describe('FactWriter.bulkWrite | 时间锚变更检测', () => {
expect(rows).toHaveLength(1);
});
});
// ─────────────────────────────────────────────
// bulkWrite — schema 违规隔离(2026-07-29 线上事故回归)
// ─────────────────────────────────────────────
/**
* 事故复盘:FRIDAY 推的结算里有 4 行「实收为负」的红冲单,manifest 当时没把它归到退费,
* 落进 payment_record 后 amount_cents<0 被 zod 拒 —— 而当时 bulkWrite 是**整批抛错**,
* ParserPipeline 因此降级成每 fact 一个事务,单批 200 行从 3 秒涨到 63.5 秒,
* 超过宿主 60s 读超时 → 对方永远收不到成功回执、无限重推同一批,卡了 13 小时一条没进。
*
* 这里锁住修复后的语义:违规的那条被隔离进 rejected,**同批其余照常写**,批不抛错。
*/
describe('FactWriter.bulkWrite | schema 违规隔离', () => {
const paymentDraft = (subjectId: string, amountCents: number): FactDraft => ({
subjectId,
kind: FactKind.ACTUAL,
type: FactType.PAYMENT_RECORD,
occurredAt: T0,
content: { payment_external_id: subjectId, amount_cents: amountCents },
});
test('一条负金额不连坐 —— 其余照常写入,违规条目进 rejected', async () => {
const { prisma, rows } = makeStore();
const { results, rejected } = await writer(prisma).bulkWrite([
{ ...baseInput, draft: paymentDraft('pay-ok-1', 12_000), transactionId: 'tx-1' },
{ ...baseInput, draft: paymentDraft('pay-bad', -50_000), transactionId: 'tx-2' },
{ ...baseInput, draft: paymentDraft('pay-ok-2', 6_000), transactionId: 'tx-3' },
]);
// 好行照常落库,批没有抛错
expect(results.map((r) => r.action)).toEqual(['created', 'created']);
expect(rows.map((r) => r.subjectId).sort()).toEqual(['pay-ok-1', 'pay-ok-2']);
// 坏行被隔离,且带够宿主定位所需的信息
expect(rejected).toHaveLength(1);
expect(rejected[0]!.subjectId).toBe('pay-bad');
expect(rejected[0]!.type).toBe(FactType.PAYMENT_RECORD);
expect(rejected[0]!.transactionId).toBe('tx-2');
expect(rejected[0]!.issues.join(' ')).toContain('amount_cents');
// 不回传 content 原文(患者信息不随回执外流)
expect(JSON.stringify(rejected[0])).not.toContain('payment_external_id');
});
test('整批全违规 → 返回空 results + 全部 rejected,仍不抛错', async () => {
const { prisma, rows } = makeStore();
const { results, rejected } = await writer(prisma).bulkWrite([
{ ...baseInput, draft: paymentDraft('bad-1', -1), transactionId: 'tx-1' },
{ ...baseInput, draft: paymentDraft('bad-2', -2), transactionId: 'tx-2' },
]);
expect(results).toHaveLength(0);
expect(rejected).toHaveLength(2);
expect(rows).toHaveLength(0);
});
test('全部合法 → rejected 为空(回归:不改动正常路径)', async () => {
const { prisma } = makeStore();
const { results, rejected } = await writer(prisma).bulkWrite([
{ ...baseInput, draft: paymentDraft('pay-1', 100), transactionId: 'tx-1' },
]);
expect(results.map((r) => r.action)).toEqual(['created']);
expect(rejected).toEqual([]);
});
});
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