Commit 0ffb6e07 by luoqi

fix(sync): 增量空转探针不再误报 —— 失败轮守卫 + 探针带业务过滤

测试服务器收到 CRITICAL「增量空转:PAC 拉到 0 行,但 DW 有新数据」。排查 sync_logs:

  07-25 16:15 UTC  failed  fetched=0  err="fatal: unexpected end of file"
  07-26 00:15 UTC  success fetched=239959 tx=103652   ← 下一轮全额补回

即:一次瞬时 ClickHouse 连接断开导致该轮失败,游标未推进,下一轮把积压全额补回
(tx=10万+,远高于平时 ~1万),**无数据丢失**。告警是误报,两处 bug 叠加:

## ① 探针在失败轮上也跑(主因)
空转告警条件只看 `tx+dup=0`,没看 status。失败轮 tx=dup=0 → 走进探针 →
把瞬时网络失败误诊成「游标格式失效/空转」。失败轮本就经日报失败计数 + error_message
暴露,不该再被探针误标。
修:加 `status===SUCCESS` 守卫,只在**成功空轮**上探。

## ② 探针反查不带业务过滤(放大误报)
探针反查只拼 `cursor > val`,而真实增量拉取还叠加业务过滤(结算正表
`is_refund=0 AND settlement_status=1`、退费表 `is_refund=1 OR settlement_status=4`、
结算方式 `settlement_status='1'`)。于是"游标之后、但会被业务条件过滤掉"的退费行
被探针数进去 → 真实拉取正确地不拉,探针却告警。
告警数据自证:settlement 正/退表各 5788 完全相等,正是这批全为退费单的铁证。
修:探针复用 extractBusinessFilters,与真实拉取同一套 WHERE。

## 验证
- 388 单测通过(新增 6 例:各表业务过滤保留/剔除、括号 OR 不被 AND 拆开)
- 失败原因「unexpected end of file」是瞬时网络错误,已自愈,无需数据补摄

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
parent 7808f53c
......@@ -182,8 +182,17 @@ export class ClickHouseSourceService {
if (!cfg.cursorValue) continue;
const sql = source.queries[tableName];
const m = sql?.match(/FROM\s+([\w.]+)/i);
if (!m) continue;
const q = `SELECT count(*) AS c FROM ${m[1]} WHERE ${cfg.cursorColumn} > '${cfg.cursorValue.replace(/'/g, "''")}'`;
if (!m || !sql) continue;
// ⭐ 探针必须与**真实增量拉取同一套 WHERE**:游标 + 业务过滤(去 cohort)。
// 只拼游标会数进"游标之后、但被业务条件过滤掉"的行(退费单 settlement_status≠1、
// 反向结算 is_refund=1、无就诊记录患者 last_visit_time NULL)→ 真实拉取正确地不拉,
// 探针却数到 → 误报"增量空转"CRITICAL(2026-07-25 实测:settlement 正/退表各 5788
// 完全相等,正是这批全为退费单的铁证)。
const clauses = [
`${cfg.cursorColumn} > '${cfg.cursorValue.replace(/'/g, "''")}'`,
...this.extractBusinessFilters(sql),
];
const q = `SELECT count(*) AS c FROM ${m[1]} WHERE ${clauses.join(' AND ')}`;
const rs = await client.query({ query: q, format: 'JSONEachRow' });
const rows = (await rs.json()) as Array<{ c: string | number }>;
out[tableName] = Number(rows[0]?.c ?? 0);
......
......@@ -1074,10 +1074,22 @@ export class ColdImportService {
// ── 增量空转探针(2026-06-10 游标格式 bug 空转三天才被发现的教训)──
// fetched=0 且确实带游标跑的增量 → 反查 DW「比游标新的行数」:
// DW 也是 0 → 真没新数据,静默;DW > 0 → PAC 拉不到但 DW 有 → 当天告警(企微)。
//
// ⭐ 只在 **status=success** 的空轮上探(2026-07-25 修):
// 失败轮(如 ClickHouse 连接中途断 "unexpected end of file")本就 tx+dup=0,且会经
// 日报失败计数 + error_message 暴露。在失败轮上跑探针,会把**瞬时网络失败误诊成
// "游标空转/格式失效"** —— 实测就发生过:一轮 failed,下一轮 tx=10万+ 全额补回(游标
// 未推进,无数据丢失),却先收到一条误导性的 CRITICAL 空转告警。
const usedCursor =
incrementalConfig &&
Object.values(incrementalConfig.perQuery).some((c) => c.cursorValue !== null);
if (!options.dryRun && usedCursor && totals.transactionsWritten + totals.duplicates === 0 && manifest.sql_source) {
if (
!options.dryRun &&
status === SyncStatus.SUCCESS &&
usedCursor &&
totals.transactionsWritten + totals.duplicates === 0 &&
manifest.sql_source
) {
const counts = await this.chSource.probeNewRowCounts(
manifest.sql_source,
incrementalConfig!.perQuery,
......
import { ClickHouseSourceService } from '../src/modules/sync/cold-import/clickhouse-source.service';
/**
* 增量空转探针回归。
*
* 探针职责:增量成功但拉 0 行时,反查 DW「比游标新的行数」,分辨
* 「DW 真没新数据(静默)」vs「游标条件失效在空转(CRITICAL 告警)」。
*
* 2026-07-25 生产误报排查暴露两个 bug:
* ① 探针反查只拼游标、**不带业务过滤**,把"游标之后但会被业务条件过滤掉"的行也数进去
* —— 退费单(settlement_status≠1 / is_refund=1)首当其冲。实测 settlement 正/退表各 5788
* 完全相等,正是这批全为退费单的铁证 → 真实拉取正确地不拉,探针却告警。
* ② 探针在 **失败轮** 上也跑(条件只看 tx+dup=0,没看 status),把瞬时网络失败
* ("unexpected end of file")误诊成"游标空转"。② 在 cold-import.service 的调用点修
* (加 status===success 守卫),① 在此处的 extractBusinessFilters 修。
*
* 本文件锁 ①:探针拼的 WHERE 必须与真实增量拉取同一套业务过滤。
*/
// extractBusinessFilters 是纯字符串逻辑、无构造依赖,直接实例化用 bracket 访问。
const svc = new ClickHouseSourceService();
const extract = (sql: string): string[] =>
(svc as unknown as { extractBusinessFilters(s: string): string[] }).extractBusinessFilters(sql);
describe('extractBusinessFilters — 探针复用真实拉取的业务过滤', () => {
test('⭐ 结算正表:保留 is_refund / settlement_status,去掉 cohort', () => {
const sql = `
SELECT id, patient_id, settlement_money
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 = extract(sql);
expect(f).toContain('settlement_status = 1');
expect(f.some((c) => /is_refund/.test(c))).toBe(true);
// cohort 子查询必须被剔除(否则探针 SQL 里嵌一坨子查询,且语义错)
expect(f.some((c) => /\(patient_id\s*,\s*brand\)\s+IN/i.test(c))).toBe(false);
});
test('⭐ 结算退费表:保留退费判定,去掉 cohort', () => {
const sql = `
SELECT id, patient_id
FROM dw_group.fact_settlement_out
WHERE (is_refund = 1 OR settlement_status = 4)
AND (patient_id, brand) IN (SELECT patient_id, brand FROM dw_group.fact_client_out WHERE last_visit_time IS NOT NULL)`;
const f = extract(sql);
expect(f.some((c) => /is_refund = 1 OR settlement_status = 4/.test(c))).toBe(true);
expect(f.some((c) => /\(patient_id\s*,\s*brand\)\s+IN/i.test(c))).toBe(false);
});
test('结算方式表:保留 settlement_status=1', () => {
const sql = `SELECT * FROM dw_group.fact_settlement_mode_out WHERE 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)`;
expect(extract(sql)).toContain("settlement_status = '1'");
});
test('主档表:last_visit_time IS NOT NULL 被剔除(增量按 cursor 拉,该条冗余)', () => {
const sql = `SELECT * FROM dw_group.fact_client_out WHERE last_visit_time IS NOT NULL`;
expect(extract(sql)).toEqual([]);
});
test('预约/病历:整段 WHERE 只有 cohort → 探针无附加过滤(与真实拉取一致)', () => {
const sql = `SELECT * FROM dw_group.fact_appointment_out
WHERE (patient_id, brand) IN (SELECT patient_id, brand FROM dw_group.fact_client_out WHERE last_visit_time IS NOT NULL)`;
expect(extract(sql)).toEqual([]);
});
test('括号内的 OR/AND 不被 AND 拆开(退费表 (a OR b) 整体保留)', () => {
const sql = `SELECT * FROM t WHERE (is_refund = 1 OR settlement_status = 4) AND settlement_status = 1`;
const f = extract(sql);
expect(f).toContain('(is_refund = 1 OR settlement_status = 4)');
expect(f).toContain('settlement_status = 1');
expect(f).toHaveLength(2);
});
});
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