Commit 4cfecec8 by luoqi

fix(sync): 增量列患者的 UNION 分支包裹 manifest query —— 别名列 + 顶层表名

测试服 dry-run 实测翻车(未落库,水位未污染):
  Missing columns: 'last_visit_time' 'patient_id' while processing query:
  'SELECT patient_id, brand FROM dw_group.fact_complex_cases_out WHERE last_visit_time > …'

listPatientPairs 构造 UNION 分支时是第三处朴素正则 `/FROM\s+([\w.]+)/`:
  ① 主档 query 现在带 SELECT 列子查询 → 抓到子查询的 FROM,表名取成 fact_complex_cases_out,
     cursor 列却仍是主档的 last_visit_time,两个错叠在一起
  ② 分支直查物理表 → 读不到 query 里的 `customer_id AS patient_id` 别名(该表物理列名是
     customer_id),patient_id 根本不存在

改为**包裹 manifest 的 query**(复用已按顶层 FROM 解析的 injectIncrementalCursor):
  SELECT <keys> FROM (SELECT … AS patient_id … FROM <表> WHERE <cursor> AND <业务过滤>)
别名在子查询里成立;顺带继承该 query 的业务过滤,与真实拉取同口径 —— 否则"游标之后但
会被业务条件过滤掉"的行会把无关患者拖进 cohort(同 extractBusinessFilters 那次的教训)。
解析不了则回退旧形态(表名改用顶层解析),不比改动前差。

补 2 项回归锁这两条,共 21 项。

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
parent 257ba23b
......@@ -449,14 +449,30 @@ export class ClickHouseSourceService {
if (incremental) {
for (const [tbl, cfg] of Object.entries(incremental.perQuery)) {
if (!cfg.cursorColumn || !cfg.cursorValue) continue;
const m = source.queries[tbl]?.match(/FROM\s+([\w.]+)/i);
const fqn = m?.[1] ?? tbl;
const q = source.queries[tbl];
const keyCols = tenant_key_column
? `${patient_key_column}, ${tenant_key_column}`
: patient_key_column;
incBranches.push(
`SELECT ${keyCols} FROM ${fqn} WHERE ${cfg.cursorColumn} > '${cfg.cursorValue.replace(/'/g, "''")}'`,
);
// ⭐ 分支 SQL **包裹 manifest 的 query**,不直接查物理表。两个原因:
// ① 患者键列名未必与 patient_key_column 同名(fact_complex_cases_out 是 customer_id,
// 靠 query 里 `AS patient_id` 对齐)——直接查表读不到别名,CH 报 Missing columns。
// ② 顺带继承该 query 的业务过滤(退费单等),与真实拉取同口径 —— 否则会因
// "游标之后但会被业务条件过滤掉"的行把无关患者拖进 cohort(同 extractBusinessFilters 的教训)。
// injectIncrementalCursor 已按顶层 FROM 解析,子查询里的 FROM/WHERE 不会串味。
// 解析不了(非常规 SQL)→ 回退旧形态:直查表名,至少不比改动前差。
let branch: string | null = null;
if (q) {
try {
branch = `SELECT ${keyCols} FROM (${this.injectIncrementalCursor(q, cfg.cursorColumn, cfg.cursorValue)})`;
} catch {
branch = null;
}
}
if (!branch) {
const fqn = (q && this.splitSelectFrom(q)?.fromTable) ?? tbl;
branch = `SELECT ${keyCols} FROM ${fqn} WHERE ${cfg.cursorColumn} > '${cfg.cursorValue.replace(/'/g, "''")}'`;
}
incBranches.push(branch);
}
}
const unionMode = incBranches.length > 0;
......
......@@ -127,6 +127,35 @@ describe('injectIncrementalCursor — 带派生列的主档 SQL', () => {
});
});
describe('增量列患者的 UNION 分支 — 别名列 + 顶层表名', () => {
/// 复刻 listPatientPairs 里 incBranches 的构造(见 clickhouse-source.service.ts)
const branchOf = (q: string, cursorCol: string, cursorVal: string) =>
`SELECT patient_id, brand FROM (${priv.injectIncrementalCursor(q, cursorCol, cursorVal)})`;
test('⭐ 主档分支查的是主档表,不是 SELECT 列里子查询的那张', () => {
const b = branchOf(MASTER_SQL, 'last_visit_time', '2026-07-30 08:15:00');
// 2026-08-01 测试服 dry-run 真实翻车过:分支被拼成
// SELECT patient_id, brand FROM dw_group.fact_complex_cases_out WHERE last_visit_time > …
// → 表名取自子查询、cursor 列取自主档,CH 报 Missing columns: 'last_visit_time' 'patient_id'
expect(b).not.toMatch(/FROM\s+dw_group\.fact_complex_cases_out\s+WHERE\s+last_visit_time/);
expect(b).toContain('FROM dw_group.fact_client_out');
});
test('⭐ 复杂病例分支能读到别名列 —— 物理表只有 customer_id,没有 patient_id', () => {
const q = `
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 b = branchOf(q, 'changed_at', '2026-07-30 08:15:00');
// 外层 SELECT patient_id 必须落在**子查询之上**(别名在那里才存在)
expect(b).toMatch(/^SELECT patient_id, brand FROM \(/);
expect(b).toContain('customer_id AS patient_id');
});
});
describe('manifest 契约', () => {
const raw = readFileSync(join(__dirname, '../data/jvs-dw/manifest.yaml'), 'utf-8');
const manifest = ColdImportManifestSchema.parse(yaml.load(raw));
......
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