Commit 0b64032d by luoqi

merge: feat/jvs-dw-complex-case-gate → test(jvs-dw 潜在治疗池不进召回 + 增量感知)

parents 5b96af15 257ba23b
Pipeline #3507 failed in 0 seconds
...@@ -26,6 +26,11 @@ field_mapping: ...@@ -26,6 +26,11 @@ field_mapping:
# 转介绍(B.1.2):推荐人数 / 带来转化总额(元;dispatcher 转分) # 转介绍(B.1.2):推荐人数 / 带来转化总额(元;dispatcher 转分)
referralCount: recommend_num referralCount: recommend_num
referralAmount: recommend_amount referralAmount: recommend_amount
# 宿主跟进闸 —— DW「潜在治疗池」(= FRIDAY 侧「复杂病例」)在跟 → PAC 不再电话召回,
# 避免两边同时联系同一个人。列由 manifest 主档 SQL 的标量子查询算出(口径/坑见那里),
# 经 canonical booleanFields coerce 成布尔 → patient_profiles.host_follow_up_active。
# 落 false(非 null)= 明确"宿主没在跟",与 FRIDAY 侧语义一致;闸判 IS NOT TRUE,行为相同。
hostFollowUpActive: has_active_complex_case
# doNotContact / deceased 不映射 — 走 PatientCanonicalSchema default false # doNotContact / deceased 不映射 — 走 PatientCanonicalSchema default false
# host(瑞尔/瑞泰)初诊一级渠道 → PAC 立柱标准(canonical-codes.PACAcquisitionChannels) # host(瑞尔/瑞泰)初诊一级渠道 → PAC 立柱标准(canonical-codes.PACAcquisitionChannels)
......
...@@ -94,6 +94,8 @@ sql_source: ...@@ -94,6 +94,8 @@ sql_source:
incremental: incremental:
per_query: per_query:
fact_client_out: { cursor_column: last_visit_time } fact_client_out: { cursor_column: last_visit_time }
# 复杂病例:用 SELECT 里的 changed_at 别名(coalesce(updated,created));CH 支持 WHERE 引用别名。
fact_complex_cases_out: { cursor_column: changed_at }
fact_appointment_out: { cursor_column: updated_date } fact_appointment_out: { cursor_column: updated_date }
fact_emr_treatment_out: { cursor_column: updated_date } fact_emr_treatment_out: { cursor_column: updated_date }
fact_settlement_out: { cursor_column: updated_date } fact_settlement_out: { cursor_column: updated_date }
...@@ -113,6 +115,15 @@ sql_source: ...@@ -113,6 +115,15 @@ sql_source:
clinic_scope: clinic_scope:
from_table: dw_group.fact_emr_treatment_out from_table: dw_group.fact_emr_treatment_out
org_column: organization_id org_column: organization_id
# 增量「反向拉主档」的来源表 —— 这些表有自己的 cursor,它们变了就要把对应患者的主档补拉回来
# (主档 cursor 是 last_visit_time,人不来诊拉不到)。前四张是历史默认(EMR 编辑等场景);
# fact_complex_cases_out 是本次新增:宿主开/关复杂病例时把人带进 cohort,跟进闸才跟得上。
reverse_pull_from:
- fact_appointment_out
- fact_emr_treatment_out
- fact_settlement_out
- fact_settlement_mode_out
- fact_complex_cases_out
# SQL 最朴素化 — host(DW)给的数据范围就是 PAC 该消化的范围。 # SQL 最朴素化 — host(DW)给的数据范围就是 PAC 该消化的范围。
# PAC 这边不预设诊所 / brand / 时间过滤,数据来什么就是什么。 # PAC 这边不预设诊所 / brand / 时间过滤,数据来什么就是什么。
...@@ -120,11 +131,39 @@ sql_source: ...@@ -120,11 +131,39 @@ sql_source:
# cohort 用 (patient_id, brand) 复合键 — patient_id 跨 brand 撞车,需带 brand 才精确。 # cohort 用 (patient_id, brand) 复合键 — patient_id 跨 brand 撞车,需带 brand 才精确。
queries: queries:
# ── 患者主档 ── # ── 患者主档 ──
# has_active_complex_case:宿主跟进闸(→ canonical hostFollowUpActive → patient_profiles)。
# DW「潜在治疗池」= FRIDAY 侧「复杂病例」:宿主自己在跟的患者,PAC 不再电话召回。
# 口径与 FRIDAY 对齐:未删除(is_del=1)且阶段 ∈ 待跟进/已咨询/已预约/诊疗中(1,2,3,5);
# 已成单(6)/已丢单(7)= 宿主停手 → 不算在跟 → 交还召回。
# ⚠️ is_del=1 是**未删除**(反直觉但两家源库一致;fact_complex_cases_out 实测 98.8% 为 1)。
# ⚠️ 复合键 (patient_id, brand):集团内同号跨品牌是两个人,单键 join 会误标瑞泰同号患者。
# ⚠️ 判定是**实时全量**的(每次拉主档都重算),不依赖增量窗口 —— 病例关闭后该列自然变 false,
# full upsert 落回 null → 患者回池。若改成"拉变化的病例再 join",增量窗口外的人会被误清。
# 注:本表只有「瑞尔」品牌数据,瑞泰患者恒 false(而非 null)—— 行为等价(闸判 IS NOT TRUE)。
fact_client_out: | fact_client_out: |
SELECT * SELECT *,
(patient_id, brand) IN (
SELECT customer_id, brand FROM dw_group.fact_complex_cases_out
WHERE is_del = 1 AND case_stage IN (1, 2, 3, 5)
) AS has_active_complex_case
FROM dw_group.fact_client_out FROM dw_group.fact_client_out
WHERE last_visit_time IS NOT NULL WHERE last_visit_time IS NOT NULL
# ── 复杂病例(潜在治疗池)──
# ⚠️ 本表**不产 fact、无 assembler**,拉它只为一件事:增量时把"病例变了但人没来诊"的患者
# 带进 cohort(见 cohort.reverse_pull_from)。真正的判定在上面主档 SQL 的标量子查询里。
# 不拉它的话:患者不来诊 → 主档 cursor(last_visit_time)拉不到 → 病例开/关 PAC 永远不知道。
# changed_at:updated_gmt_at 有 44% 为 NULL(建后没改过,恰是刚入池的新病例),
# 直接拿它当 cursor 会 `NULL > x` 恒 UNKNOWN 静默漏掉这批 → 必须 coalesce 到 created_gmt_at。
fact_complex_cases_out: |
SELECT customer_id AS patient_id, brand,
coalesce(updated_gmt_at, toDateTime(created_gmt_at)) AS changed_at
FROM dw_group.fact_complex_cases_out
WHERE (customer_id, brand) IN (
SELECT patient_id, brand FROM dw_group.fact_client_out
WHERE last_visit_time IS NOT NULL
)
# ── 预约(全部状态;in_time NULL 的在 transforms.route 时分流跳过 encounter) ── # ── 预约(全部状态;in_time NULL 的在 transforms.route 时分流跳过 encounter) ──
fact_appointment_out: | fact_appointment_out: |
SELECT * SELECT *
......
import { Injectable, Logger } from '@nestjs/common'; import { Injectable, Logger } from '@nestjs/common';
import * as fs from 'node:fs'; import * as fs from 'node:fs';
import { createClient, type ClickHouseClient } from '@clickhouse/client'; import { createClient, type ClickHouseClient } from '@clickhouse/client';
import type { ClickHouseSource } from './manifest.schema'; import { DEFAULT_REVERSE_PULL_FROM, type ClickHouseSource } from './manifest.schema';
/** /**
* 解析 PAC_COHORT_ONLY_PATIENT —— 定向重摄的患者 id 集合(cohort 收窄用)。 * 解析 PAC_COHORT_ONLY_PATIENT —— 定向重摄的患者 id 集合(cohort 收窄用)。
...@@ -259,7 +259,7 @@ export class ClickHouseSourceService { ...@@ -259,7 +259,7 @@ export class ClickHouseSourceService {
// 合并入已有 fact_client_out tables(set 去重)→ 数据完整性 100% // 合并入已有 fact_client_out tables(set 去重)→ 数据完整性 100%
// stub 兜底(方案 A,parser 侧)+ 主档反向拉(方案 C,sync 侧)= 完整双保险 // stub 兜底(方案 A,parser 侧)+ 主档反向拉(方案 C,sync 侧)= 完整双保险
if (incremental && tables['fact_client_out']) { if (incremental && tables['fact_client_out']) {
const reverseRows = await this.reversePullPatientMaster(client, tables); const reverseRows = await this.reversePullPatientMaster(client, tables, source);
if (reverseRows.length > 0) { if (reverseRows.length > 0) {
const before = tables['fact_client_out'].length; const before = tables['fact_client_out'].length;
// 去重 by patient_id+brand:已有 cursor 拉到的不重复 push // 去重 by patient_id+brand:已有 cursor 拉到的不重复 push
...@@ -294,34 +294,51 @@ export class ClickHouseSourceService { ...@@ -294,34 +294,51 @@ export class ClickHouseSourceService {
private async reversePullPatientMaster( private async reversePullPatientMaster(
client: ClickHouseClient, client: ClickHouseClient,
tables: Record<string, unknown[]>, tables: Record<string, unknown[]>,
source: ClickHouseSource,
): Promise<unknown[]> { ): Promise<unknown[]> {
// 收集所有 fact 表的 patient_id + brand 集合(去重) // 来源表 / 键列 / 主档表全部从 manifest 读(宿主口径),不配则回落历史默认。
const factTables = ['fact_appointment_out', 'fact_emr_treatment_out', 'fact_settlement_out', 'fact_settlement_mode_out']; // 此前四张 jvs-dw 表名硬编码在此 —— 新增一张要改通用代码,与"yaml 是宿主唯一差异"相悖。
const cohort = source.cohort;
const factTables = cohort?.reverse_pull_from ?? [...DEFAULT_REVERSE_PULL_FROM];
const keyCol = cohort?.patient_key_column ?? 'patient_id';
const tenantCol = cohort?.tenant_key_column;
const masterTable = cohort?.patient_list_from ?? 'dw_group.fact_client_out';
// 收集来源表的患者键集合(去重)。配了 tenant_key_column → 复合键(集团内同号跨品牌是两个人)
const pairs = new Set<string>(); const pairs = new Set<string>();
for (const t of factTables) { for (const t of factTables) {
for (const row of (tables[t] ?? []) as Record<string, unknown>[]) { for (const row of (tables[t] ?? []) as Record<string, unknown>[]) {
const pid = row.patient_id; const pid = row[keyCol];
const brand = row.brand; const brand = tenantCol ? row[tenantCol] : undefined;
if (pid !== null && pid !== undefined && brand) { if (pid !== null && pid !== undefined && (!tenantCol || brand)) {
pairs.add(`${pid}|||${brand}`); pairs.add(`${pid}|||${brand ?? ''}`);
} }
} }
} }
if (pairs.size === 0) return []; if (pairs.size === 0) return [];
// 构造 IN ((pid1,brand1),(pid2,brand2),...) tuple list // 构造 IN ((pid1,brand1),(pid2,brand2),...) tuple list
// CH 支持 `(patient_id, brand) IN ((1,'瑞尔'),(2,'瑞泰'))` tuple in 语法 // CH 支持 `(patient_id, brand) IN ((1,'瑞尔'),(2,'瑞泰'))` tuple in 语法
const esc = (v: string) => v.replace(/'/g, "''");
const tuples = [...pairs] const tuples = [...pairs]
.map((p) => { .map((p) => {
const [pid, brand] = p.split('|||'); const [pid, brand] = p.split('|||');
return `('${(pid ?? '').replace(/'/g, "''")}', '${(brand ?? '').replace(/'/g, "''")}')`; return tenantCol
? `('${esc(pid ?? '')}', '${esc(brand ?? '')}')`
: `'${esc(pid ?? '')}'`;
}) })
.join(','); .join(',');
const sql = `SELECT * FROM dw_group.fact_client_out WHERE (patient_id, brand) IN (${tuples})`; const keyExpr = tenantCol ? `(${keyCol}, ${tenantCol})` : keyCol;
// ⭐ 复用主档 query 的 SELECT 列,**不能写死 `SELECT *`**:
// 主档 SQL 可能带派生列(jvs-dw 的 has_active_complex_case 由标量子查询算出),
// `SELECT *` 拉回来就少这一列 → canonical 无此键 → full 分支落 null → 正在跟进的患者
// 被反向拉这一步静默"洗回"召回池。取不到原 SQL 时回退 `*`(行为同旧版)。
const masterQuery = source.queries?.[this.tableKeyOf(masterTable)];
const cols = (masterQuery && this.splitSelectFrom(masterQuery)?.selectCols) || '*';
const sql = `SELECT ${cols} FROM ${masterTable} WHERE ${keyExpr} IN (${tuples})`;
this.logger.log(`[clickhouse] 反向拉主档 query ${pairs.size} pairs`); this.logger.log(`[clickhouse] 反向拉主档 query ${pairs.size} pairs`);
const started = Date.now(); const started = Date.now();
const result = await client.query({ query: sql, format: 'JSONEachRow' }); const result = await client.query({ query: sql, format: 'JSONEachRow' });
const rows = (await result.json()) as unknown[]; const rows = (await result.json()) as unknown[];
this.logger.log(`[clickhouse] 反向 fact_client_out → ${rows.length} 行,${Date.now() - started} ms`); this.logger.log(`[clickhouse] 反向 ${masterTable} ${rows.length} ,${Date.now() - started} ms`);
return rows; return rows;
} }
...@@ -678,14 +695,13 @@ export class ClickHouseSourceService { ...@@ -678,14 +695,13 @@ export class ClickHouseSourceService {
cursorValue: string | null, cursorValue: string | null,
): string { ): string {
// 提取 SELECT 子句 (FROM 之前) 跟 FROM <table> // 提取 SELECT 子句 (FROM 之前) 跟 FROM <table>
const m = originalSql.match(/^\s*SELECT\s+([\s\S]+?)\s+FROM\s+([\w.]+)/i); const parsed = this.splitSelectFrom(originalSql);
if (!m) { if (!parsed) {
throw new Error( throw new Error(
`[incremental] cannot parse SQL for cursor injection: ${originalSql.slice(0, 80)}...`, `[incremental] cannot parse SQL for cursor injection: ${originalSql.slice(0, 80)}...`,
); );
} }
const selectCols = m[1]!.trim(); const { selectCols, fromTable } = parsed;
const fromTable = m[2]!;
// 收集所有 WHERE 子句(cursor + 业务过滤),最后统一 'WHERE ... AND ...' 拼装 // 收集所有 WHERE 子句(cursor + 业务过滤),最后统一 'WHERE ... AND ...' 拼装
const clauses: string[] = []; const clauses: string[] = [];
if (cursorValue) { if (cursorValue) {
...@@ -700,10 +716,17 @@ export class ClickHouseSourceService { ...@@ -700,10 +716,17 @@ export class ClickHouseSourceService {
/// 从原 SQL 抽出 WHERE 子句里**非 cohort** 的业务过滤条件 /// 从原 SQL 抽出 WHERE 子句里**非 cohort** 的业务过滤条件
/// 返回 clause 数组(每个 clause 不带 WHERE/AND 前缀,纯条件) /// 返回 clause 数组(每个 clause 不带 WHERE/AND 前缀,纯条件)
private extractBusinessFilters(sql: string): string[] { private extractBusinessFilters(sql: string): string[] {
// 抓 WHERE ... 到 ORDER BY/LIMIT/end 之间 // 抓**顶层** WHERE ... 到 ORDER BY/LIMIT/end 之间。
const m = sql.match(/\bWHERE\s+([\s\S]+?)(?:\bORDER\s+BY\b|\bLIMIT\b|$)/i); // ⚠️ 必须找顶层:SELECT 列里的标量子查询也带 WHERE(如 jvs-dw 主档的
if (!m) return []; // `(patient_id,brand) IN (SELECT … WHERE is_del=1 …) AS has_active_complex_case`),
const whereBody = m[1]!.trim(); // 用 /\bWHERE\b/ 抓第一个会把子查询的过滤当成本表的业务过滤搬到外层 → 列不存在,CH 报错。
const whereAt = this.topLevelIndexOf(sql, 'WHERE');
if (whereAt < 0) return [];
const tail = sql.slice(whereAt + 'WHERE'.length);
const stopAt = [this.topLevelIndexOf(tail, 'ORDER'), this.topLevelIndexOf(tail, 'LIMIT')]
.filter((i) => i >= 0)
.reduce((min, i) => (min < 0 || i < min ? i : min), -1);
const whereBody = (stopAt >= 0 ? tail.slice(0, stopAt) : tail).trim();
// 拆 AND,去掉:① cohort 子查询 ② 老的 last_visit_time IS NOT NULL(增量不需要) // 拆 AND,去掉:① cohort 子查询 ② 老的 last_visit_time IS NOT NULL(增量不需要)
return this.splitAndClauses(whereBody) return this.splitAndClauses(whereBody)
.map((c) => c.trim()) .map((c) => c.trim())
...@@ -715,6 +738,50 @@ export class ClickHouseSourceService { ...@@ -715,6 +738,50 @@ export class ClickHouseSourceService {
); );
} }
/// 表全名 → queries 的 key(queries 用不带库前缀的表名,manifest cohort 里写的是全名)
private tableKeyOf(fullName: string): string {
const i = fullName.lastIndexOf('.');
return i >= 0 ? fullName.slice(i + 1) : fullName;
}
/**
* 找**顶层**(括号深度 0)关键字的下标;找不到返回 -1。大小写不敏感,按词边界匹配。
*
* 为什么需要:SQL 改写用的正则都假设"第一个 FROM / WHERE 就是本表的",
* 一旦 SELECT 列里出现标量子查询(`… IN (SELECT … FROM … WHERE …) AS flag`)就全错位。
* 括号深度是区分"我的子句"和"子查询的子句"的唯一可靠依据 —— 与 splitAndClauses 同款做法。
*/
private topLevelIndexOf(sql: string, keyword: string): number {
const re = new RegExp(`\\b${keyword}\\b`, 'gi');
let depth = 0;
let scanned = 0;
let m: RegExpExecArray | null;
while ((m = re.exec(sql)) !== null) {
for (; scanned < m.index; scanned++) {
const ch = sql[scanned];
if (ch === '(') depth++;
else if (ch === ')') depth--;
}
if (depth === 0) return m.index;
}
return -1;
}
/**
* 拆 `SELECT <cols> FROM <table>` —— 按顶层 FROM 切,子查询里的 FROM 不参与。
* 返回 null = 不是可识别的单表 SELECT(调用方各自决定报错还是回退)。
*/
private splitSelectFrom(sql: string): { selectCols: string; fromTable: string } | null {
if (!/^\s*SELECT\b/i.test(sql)) return null;
const fromAt = this.topLevelIndexOf(sql, 'FROM');
if (fromAt < 0) return null;
const selectCols = sql.slice(sql.search(/\bSELECT\b/i) + 'SELECT'.length, fromAt).trim();
const rest = sql.slice(fromAt + 'FROM'.length).trim();
const tbl = rest.match(/^([\w.]+)/);
if (!selectCols || !tbl) return null;
return { selectCols, fromTable: tbl[1]! };
}
private splitAndClauses(whereBody: string): string[] { private splitAndClauses(whereBody: string): string[] {
// 按 ` AND `(大小写)拆分,但跳过括号内的 AND // 按 ` AND `(大小写)拆分,但跳过括号内的 AND
const out: string[] = []; const out: string[] = [];
......
...@@ -120,9 +120,31 @@ export const ClickHouseSourceSchema = z.object({ ...@@ -120,9 +120,31 @@ export const ClickHouseSourceSchema = z.object({
org_column: z.string().min(1).default('organization_id'), org_column: z.string().min(1).default('organization_id'),
}) })
.optional(), .optional(),
/// 增量「反向拉主档」的来源表(queries 里的 key,不带库前缀)。
///
/// 【解决什么】主档 cursor 是 list_cursor_column(jvs-dw = last_visit_time),
/// 患者不来诊则主档拉不到他 —— 但他的**事实**可能变了(EMR 被编辑、复杂病例被关闭)。
/// 增量跑完后从这些表里收集患者键,反向补拉一次主档,保证"事实变了主档也在场"。
///
/// 【为什么要可配】此前是四张 jvs-dw 表名硬编码在 service 里,新增一张要改通用代码;
/// 而"哪些表的变化该带出主档"本就是**宿主口径**(取决于该宿主哪些表有独立 cursor)。
/// 不配 → 保持历史默认(见 DEFAULT_REVERSE_PULL_FROM),行为不变。
///
/// ⚠️ 表的行必须带 patient_key_column / tenant_key_column 两列(SELECT 里 AS 别名亦可),
/// 否则收集不到键、该表静默不参与反向拉。
reverse_pull_from: z.array(z.string().min(1)).optional(),
}) })
.optional(), .optional(),
}); });
/// 反向拉主档的历史默认表集(jvs-dw 首版硬编码,保留为默认值以免既有 manifest 行为变化)。
/// 新宿主 / 新增表请在 manifest 的 cohort.reverse_pull_from 显式声明。
export const DEFAULT_REVERSE_PULL_FROM = [
'fact_appointment_out',
'fact_emr_treatment_out',
'fact_settlement_out',
'fact_settlement_mode_out',
] as const;
export type ClickHouseSource = z.infer<typeof ClickHouseSourceSchema>; export type ClickHouseSource = z.infer<typeof ClickHouseSourceSchema>;
export const ColdImportManifestSchema = z export const ColdImportManifestSchema = z
......
import { readFileSync } from 'node:fs';
import { join } from 'node:path';
import * as yaml from 'js-yaml';
import { ClickHouseSourceService } from '../src/modules/sync/cold-import/clickhouse-source.service';
import {
ColdImportManifestSchema,
DEFAULT_REVERSE_PULL_FROM,
} from '../src/modules/sync/cold-import/manifest.schema';
/**
* jvs-dw 宿主跟进闸(DW「潜在治疗池」= FRIDAY 侧「复杂病例」)的 SQL 改写回归。
*
* 【为什么这组测试必须存在】
* 接这个字段时,主档 query 头一次出现**SELECT 列里的标量子查询**:
* SELECT *, (patient_id,brand) IN (SELECT … FROM … WHERE is_del=1 …) AS has_active_complex_case
* FROM dw_group.fact_client_out WHERE last_visit_time IS NOT NULL
* 而增量链路的两个 SQL 改写器原本都假设"第一个 FROM / WHERE 就是本表的":
* - injectIncrementalCursor:/SELECT ([\s\S]+?) FROM ([\w.]+)/ 非贪婪 → 撞上子查询的 FROM,
* 切出半截列表 + 错误表名,增量直接拉错表
* - extractBusinessFilters:/\bWHERE\b/ 抓第一个 → 把子查询的 is_del=1 当成主档的业务过滤
* 搬到外层 → 主档没有 is_del 列 → CH 报错(且探针 SQL 一并错)
* 两者都改成按**括号深度**定位顶层关键字。这些是静默错法(不崩就是拉错数据),必须锁死。
*
* 另锁两条与"闸"本身的正确性直接相关的纪律:
* - 反向拉主档必须复用主档 query 的 SELECT 列(写死 `SELECT *` 会丢派生列 → 正在跟进的
* 患者被 full upsert 洗回 null → 静默回到召回池)
* - manifest 真文件能过 schema,且 reverse_pull_from 确实带上了复杂病例表
*/
const svc = new ClickHouseSourceService();
type Priv = {
extractBusinessFilters(s: string): string[];
splitSelectFrom(s: string): { selectCols: string; fromTable: string } | null;
topLevelIndexOf(s: string, kw: string): number;
injectIncrementalCursor(s: string, col: string, val: string | null): string;
tableKeyOf(s: string): string;
};
const priv = svc as unknown as Priv;
/// 与 data/jvs-dw/manifest.yaml 的主档 query 同形(带派生列的标量子查询)
const MASTER_SQL = `
SELECT *,
(patient_id, brand) IN (
SELECT customer_id, brand FROM dw_group.fact_complex_cases_out
WHERE is_del = 1 AND case_stage IN (1, 2, 3, 5)
) AS has_active_complex_case
FROM dw_group.fact_client_out
WHERE last_visit_time IS NOT NULL`;
describe('顶层关键字定位 — 子查询里的 FROM/WHERE 不参与', () => {
test('splitSelectFrom 取顶层 FROM,不是子查询的', () => {
const parsed = priv.splitSelectFrom(MASTER_SQL);
expect(parsed).not.toBeNull();
// ⭐ 关键:表名是主档,不是子查询里的 fact_complex_cases_out
expect(parsed!.fromTable).toBe('dw_group.fact_client_out');
// 派生列完整保留(含整个子查询和别名),否则增量拉回来就少这一列
expect(parsed!.selectCols).toContain('has_active_complex_case');
expect(parsed!.selectCols).toContain('fact_complex_cases_out');
// 括号必须配平 —— 半截子查询是老正则最典型的产物
const open = (parsed!.selectCols.match(/\(/g) ?? []).length;
const close = (parsed!.selectCols.match(/\)/g) ?? []).length;
expect(open).toBe(close);
});
test('topLevelIndexOf 跳过括号内的同名关键字', () => {
const at = priv.topLevelIndexOf(MASTER_SQL, 'WHERE');
// 顶层 WHERE 后面跟的是 last_visit_time,不是子查询的 is_del
expect(MASTER_SQL.slice(at)).toMatch(/^WHERE\s+last_visit_time/);
});
test('无顶层 FROM / 非 SELECT → null(调用方自行报错,不硬崩)', () => {
expect(priv.splitSelectFrom('SHOW TABLES')).toBeNull();
expect(priv.splitSelectFrom('SELECT 1')).toBeNull();
});
});
describe('extractBusinessFilters — 不把子查询的过滤搬到外层', () => {
test('⭐ 主档 SQL:只留 last_visit_time,绝不带出 is_del / case_stage', () => {
const f = priv.extractBusinessFilters(MASTER_SQL);
// is_del / case_stage 是**子查询**里的条件,搬到主档外层会因列不存在直接报错
expect(f.some((c) => /is_del/i.test(c))).toBe(false);
expect(f.some((c) => /case_stage/i.test(c))).toBe(false);
// last_visit_time IS NOT NULL 属于被显式剔除的一类(增量不需要)→ 结果应为空
expect(f).toHaveLength(0);
});
test('回归:结算表的业务过滤照旧保留(老行为不变)', () => {
const sql = `
SELECT id, patient_id FROM dw_group.fact_settlement_out
WHERE (is_refund IS NULL OR is_refund = 0)
AND settlement_status = 1
AND (patient_id, brand) IN (
SELECT patient_id, brand FROM dw_group.fact_client_out WHERE last_visit_time IS NOT NULL
)`;
const f = priv.extractBusinessFilters(sql);
expect(f).toContain('settlement_status = 1');
expect(f.some((c) => /\(patient_id\s*,\s*brand\)\s+IN/i.test(c))).toBe(false);
});
});
describe('injectIncrementalCursor — 带派生列的主档 SQL', () => {
test('⭐ 表名正确 + 派生列保留 + 子查询条件不外泄', () => {
const out = priv.injectIncrementalCursor(MASTER_SQL, 'last_visit_time', '2026-07-01 00:00:00');
expect(out).toContain('FROM dw_group.fact_client_out');
expect(out).not.toMatch(/FROM\s+dw_group\.fact_complex_cases_out\s+WHERE\s+last_visit_time/);
expect(out).toContain('has_active_complex_case');
expect(out).toContain("last_visit_time > '2026-07-01 00:00:00'");
// 子查询的条件不能出现在**顶层** WHERE(它仍应留在 SELECT 列的括号里)
const topWhere = out.slice(priv.topLevelIndexOf(out, 'WHERE'));
expect(topWhere).not.toMatch(/\bis_del\b/);
});
test('复杂病例表:cursor 用 SELECT 别名 changed_at(源列 44% NULL 不能直接当 cursor)', () => {
const sql = `
SELECT customer_id AS patient_id, brand,
coalesce(updated_gmt_at, toDateTime(created_gmt_at)) AS changed_at
FROM dw_group.fact_complex_cases_out
WHERE (customer_id, brand) IN (
SELECT patient_id, brand FROM dw_group.fact_client_out WHERE last_visit_time IS NOT NULL
)`;
const out = priv.injectIncrementalCursor(sql, 'changed_at', '2026-07-01 00:00:00');
expect(out).toContain('FROM dw_group.fact_complex_cases_out');
// 别名定义必须还在(WHERE 引用 SELECT 别名是 CH 的能力,别名被切掉就失效)
expect(out).toContain('AS changed_at');
expect(out).toContain("changed_at > '2026-07-01 00:00:00'");
expect(out).toContain('ORDER BY changed_at');
});
});
describe('manifest 契约', () => {
const raw = readFileSync(join(__dirname, '../data/jvs-dw/manifest.yaml'), 'utf-8');
const manifest = ColdImportManifestSchema.parse(yaml.load(raw));
test('jvs-dw manifest 过 schema', () => {
expect(manifest.host_name).toBe('jvs-dw');
});
test('⭐ reverse_pull_from 带上复杂病例表 —— 否则病例开/关时人不来诊就永远感知不到', () => {
const list = manifest.sql_source?.cohort?.reverse_pull_from ?? [];
expect(list).toContain('fact_complex_cases_out');
// 历史四张不能丢(丢了 = EMR 编辑等既有场景静默退化)
for (const t of DEFAULT_REVERSE_PULL_FROM) expect(list).toContain(t);
});
test('复杂病例表既在 queries 里、也配了 cursor(否则增量拉不到变化)', () => {
expect(manifest.sql_source?.queries?.['fact_complex_cases_out']).toBeDefined();
expect(
manifest.sql_source?.incremental?.per_query?.['fact_complex_cases_out']?.cursor_column,
).toBe('changed_at');
});
test('主档 query 产出 has_active_complex_case,且用复合键 (patient_id, brand)', () => {
const sql = manifest.sql_source?.queries?.['fact_client_out'] ?? '';
expect(sql).toContain('has_active_complex_case');
// ⚠️ 单键 join 会把瑞泰同号患者误标成"在跟"(集团内同号跨品牌是两个人)
expect(sql).toMatch(/\(patient_id,\s*brand\)\s+IN/);
// 口径:未删除 + 在跟四阶段
expect(sql).toMatch(/is_del\s*=\s*1/);
expect(sql).toMatch(/case_stage\s+IN\s*\(1,\s*2,\s*3,\s*5\)/);
});
test('patient assembler 映射到 canonical hostFollowUpActive', () => {
const y = readFileSync(join(__dirname, '../data/jvs-dw/assemblers/patient.yaml'), 'utf-8');
const cfg = yaml.load(y) as { field_mapping: Record<string, string> };
expect(cfg.field_mapping.hostFollowUpActive).toBe('has_active_complex_case');
});
test('tableKeyOf 去库前缀(cohort 写全名,queries 用短名)', () => {
expect(priv.tableKeyOf('dw_group.fact_client_out')).toBe('fact_client_out');
expect(priv.tableKeyOf('fact_client_out')).toBe('fact_client_out');
});
});
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