Commit 2359bbf7 by luoqi

feat(sync): cold-import 新增 --since/--months 按时间收窄存量摄入 cohort

和 --clinics 同套路:只收窄"列哪些患者"(按 cohort.list_cursor_column=last_visit_time),
被选中患者的全部历史仍全摄;可与 --clinics 叠加(AND)。不用改 manifest。
  --since=2026-01-01   今年有来诊的患者
  --months=12         最近 12 个月有来诊(CLI 换算成 cutoff 日期,--since 显式优先)
DW 实测:全量 277万 → 今年 46万 / 近12月 75万 / 近24月 118万。
last_visit_time 是 ISO 文本,字典序=时间序,直接串比较(避 NO_COMMON_TYPE)。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
parent fcd77227
Pipeline #3367 failed in 0 seconds
...@@ -23,6 +23,8 @@ interface CliArgs { ...@@ -23,6 +23,8 @@ interface CliArgs {
incremental: boolean; incremental: boolean;
cohortBatchSize?: number | null; // null = 显式禁用分批,undefined = 用 env / 默认 cohortBatchSize?: number | null; // null = 显式禁用分批,undefined = 用 env / 默认
clinics?: string[]; // --clinics=X,Y:只列"在这些诊所看过"的患者(其全部相关主体仍全摄);不传=全量 clinics?: string[]; // --clinics=X,Y:只列"在这些诊所看过"的患者(其全部相关主体仍全摄);不传=全量
since?: string; // 收窄到"该日期后有来诊"的患者(YYYY-MM-DD);--months=N 会换算成日期
months?: number; // --months=N:最近 N 个月(换算成 since 的 cutoff 日期)
} }
function parseArgs(argv: string[]): CliArgs { function parseArgs(argv: string[]): CliArgs {
...@@ -41,8 +43,19 @@ function parseArgs(argv: string[]): CliArgs { ...@@ -41,8 +43,19 @@ function parseArgs(argv: string[]): CliArgs {
.split(',') .split(',')
.map((s) => s.trim()) .map((s) => s.trim())
.filter(Boolean); .filter(Boolean);
} else if (a.startsWith('--since=')) {
args.since = a.slice('--since='.length).trim();
} else if (a.startsWith('--months=')) {
const n = parseInt(a.slice('--months='.length), 10);
if (Number.isFinite(n) && n > 0) args.months = n;
} else if (a.startsWith('--dir=')) args.dir = a.slice('--dir='.length); } else if (a.startsWith('--dir=')) args.dir = a.slice('--dir='.length);
} }
// --since 显式优先;否则 --months=N 换算成 cutoff 日期(今天往前推 N 个月)。
if (!args.since && args.months && args.months > 0) {
const d = new Date();
d.setMonth(d.getMonth() - args.months);
args.since = d.toISOString().slice(0, 10); // YYYY-MM-DD
}
return args; return args;
} }
...@@ -64,12 +77,19 @@ function printHelp() { ...@@ -64,12 +77,19 @@ function printHelp() {
' --no-cohort 显式禁用分批,跑 single-shot(文件源或调试场景)', ' --no-cohort 显式禁用分批,跑 single-shot(文件源或调试场景)',
' --clinics=<id,id> 只列"在这些诊所看过"的患者(其全部相关主体仍全摄;需 manifest', ' --clinics=<id,id> 只列"在这些诊所看过"的患者(其全部相关主体仍全摄;需 manifest',
' cohort.clinic_scope 配置)。不传=全量。用于按诊所分批补摄,幂等追加不清空。', ' cohort.clinic_scope 配置)。不传=全量。用于按诊所分批补摄,幂等追加不清空。',
' --since=<YYYY-MM-DD> 只列"该日期后有来诊"的患者(按 cohort.list_cursor_column 收窄;',
' 被选中患者的全部历史仍全摄)。如 --since=2026-01-01 = 今年。',
' --months=<N> 只列"最近 N 个月有来诊"的患者(换算成 --since 的 cutoff;--since 显式优先)。',
' --clinics 与 --since/--months 可叠加(AND 求交)。',
' --help, -h 显示本帮助', ' --help, -h 显示本帮助',
'', '',
'Examples:', 'Examples:',
' pnpm cold-import -- --dir=./data/jvs-dw --dry-run', ' pnpm cold-import -- --dir=./data/jvs-dw --dry-run',
' pnpm cold-import -- --dir=./data/jvs-dw', ' pnpm cold-import -- --dir=./data/jvs-dw',
' pnpm cold-import -- --dir=./data/jvs-dw --cohort-batch=2000', ' pnpm cold-import -- --dir=./data/jvs-dw --cohort-batch=2000',
' pnpm cold-import -- --dir=./data/jvs-dw --months=12 # 最近 12 个月有来诊的患者',
' pnpm cold-import -- --dir=./data/jvs-dw --since=2026-01-01 # 今年有来诊的患者',
' pnpm cold-import -- --dir=./data/jvs-dw --clinics=<orgId> --months=24 # 某诊所近 24 月(叠加)',
' pnpm cold-import -- --dir=./data/jvs-dw --incremental', ' pnpm cold-import -- --dir=./data/jvs-dw --incremental',
'', '',
]; ];
...@@ -87,7 +107,9 @@ async function bootstrap() { ...@@ -87,7 +107,9 @@ async function bootstrap() {
const logger = new Logger('cold-import:cli'); const logger = new Logger('cold-import:cli');
logger.log( logger.log(
`Starting cold-import CLI(dir=${args.dir}, dryRun=${args.dryRun}, ` + `Starting cold-import CLI(dir=${args.dir}, dryRun=${args.dryRun}, ` +
`incremental=${args.incremental}, cohortBatch=${args.cohortBatchSize ?? '(default/env)'})`, `incremental=${args.incremental}, cohortBatch=${args.cohortBatchSize ?? '(default/env)'}` +
`${args.clinics?.length ? `, clinics=${args.clinics.join(',')}` : ''}` +
`${args.since ? `, since=${args.since}${args.months ? `(--months=${args.months})` : ''}` : ''})`,
); );
const app = await NestFactory.createApplicationContext(AppModule, { const app = await NestFactory.createApplicationContext(AppModule, {
...@@ -101,6 +123,7 @@ async function bootstrap() { ...@@ -101,6 +123,7 @@ async function bootstrap() {
incremental: args.incremental, incremental: args.incremental,
cohortBatchSize: args.cohortBatchSize, cohortBatchSize: args.cohortBatchSize,
clinics: args.clinics, clinics: args.clinics,
since: args.since,
}); });
logger.log('─────────────────────────────────────────'); logger.log('─────────────────────────────────────────');
......
...@@ -332,6 +332,7 @@ export class ClickHouseSourceService { ...@@ -332,6 +332,7 @@ export class ClickHouseSourceService {
source: ClickHouseSource, source: ClickHouseSource,
incremental?: IncrementalConfig, incremental?: IncrementalConfig,
clinics?: string[], clinics?: string[],
since?: string,
): Promise<CohortKey[]> { ): Promise<CohortKey[]> {
const cohort = source.cohort; const cohort = source.cohort;
if (!cohort) { if (!cohort) {
...@@ -413,6 +414,20 @@ export class ClickHouseSourceService { ...@@ -413,6 +414,20 @@ export class ClickHouseSourceService {
); );
this.logger.log(`[clickhouse·cohort] --clinics 收窄:${clinics.length} 家诊所 (${cs.from_table}.${cs.org_column})`); this.logger.log(`[clickhouse·cohort] --clinics 收窄:${clinics.length} 家诊所 (${cs.from_table}.${cs.org_column})`);
} }
// --since=<date> / --months=N:把 cohort 收窄到「该日期后有来诊」的患者,按 list_cursor_column
// (患者主档时间列 last_visit_time)。同样只收窄「列哪些患者」——被选中患者的全部历史
// (含更早的事件)照常摄入(不误删)。可与 --clinics 叠加(AND)。需 cohort.list_cursor_column。
if (since) {
if (!list_cursor_column) {
throw new Error(
'--since/--months 需要 manifest.sql_source.cohort.list_cursor_column 配置(患者主档时间列)',
);
}
// DW 的 last_visit_time 是 String(ISO 文本):空串 < 任何日期,故 >= cutoff 天然排除空串;
// ISO 日期字典序 = 时间序,直接串比较,不转类型(避 NO_COMMON_TYPE)。
whereParts.push(`${list_cursor_column} >= '${since.replace(/'/g, "''")}'`);
this.logger.log(`[clickhouse·cohort] --since 收窄:${list_cursor_column} >= ${since}`);
}
const whereSql = whereParts.length > 0 ? ` WHERE ${whereParts.join(' AND ')}` : ''; const whereSql = whereParts.length > 0 ? ` WHERE ${whereParts.join(' AND ')}` : '';
// dev/ops:PAC_COHORT_LIMIT=N → 只取 N 个患者(本地抽样重摄用)。不设 = 全量,默认行为不变。 // dev/ops:PAC_COHORT_LIMIT=N → 只取 N 个患者(本地抽样重摄用)。不设 = 全量,默认行为不变。
// PAC_COHORT_SAMPLE = recent(默认,最近就诊优先)| oldest(最久未来,lapsed,召回候选多) // PAC_COHORT_SAMPLE = recent(默认,最近就诊优先)| oldest(最久未来,lapsed,召回候选多)
......
...@@ -587,6 +587,9 @@ export class ColdImportService { ...@@ -587,6 +587,9 @@ export class ColdImportService {
/// --clinics=X,Y:只列"在这些诊所看过"的患者(收窄 cohort;其全部相关主体仍全摄)。 /// --clinics=X,Y:只列"在这些诊所看过"的患者(收窄 cohort;其全部相关主体仍全摄)。
/// 需 manifest cohort.clinic_scope 配置。不传/空 → 全量(默认行为不变)。 /// 需 manifest cohort.clinic_scope 配置。不传/空 → 全量(默认行为不变)。
clinics?: string[]; clinics?: string[];
/// --since=<YYYY-MM-DD>:只列"该日期后有来诊"的患者(按 cohort.list_cursor_column 收窄;
/// 被选中患者的全部历史仍全摄)。CLI 的 --months=N 会换算成日期传进来。可与 clinics 叠加。
since?: string;
} = {}, } = {},
): Promise<ImportRunResult> { ): Promise<ImportRunResult> {
const absDir = path.resolve(dir); const absDir = path.resolve(dir);
...@@ -765,6 +768,7 @@ export class ColdImportService { ...@@ -765,6 +768,7 @@ export class ColdImportService {
sqlSource, sqlSource,
incrementalConfig, incrementalConfig,
options.clinics, options.clinics,
options.since,
); );
if (allPairs.length === 0) { if (allPairs.length === 0) {
// A:增量模式 + cohort 空 → DW 本窗口(已含回看)无任何数据 → 标记不推进游标 // A:增量模式 + cohort 空 → DW 本窗口(已含回看)无任何数据 → 标记不推进游标
......
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