Commit 9059386b by luoqi

refactor(sync): 增量水位声明与源类型解耦 — file 源存量跑同产 cursor(与拉模式同一段逻辑)

统一动机:jvs-dw(拉)存量跑写 cursor_after=run_start 供首次增量接力;FRIDAY(file)
此前不写 —— 差异根源不是流程分叉(finally 记账本就单点共用),而是 incremental.per_query
声明寄生在 sql_source 之下,file 源无处声明水位列。

- manifest.schema:IncrementalDeclSchema 抽名;sql_source.incremental 引用之(形状不动),
  manifest 顶层新增可选 incremental(file 源用;两处都写时 sql_source 内优先)
- cold-import:入口一行 ?? 兜底(sql_source 优先短路)——之后读写 cursor 走完全同一段
  finally,零新逻辑。file 装载路径不消费水位(loadAllTables 文件分支忽略 incremental
  参数,全量装载幂等去重),水位纯记账,供 delta 导出 WHERE 模板 / push 回放起点
- friday manifest:顶层声明 12 张事件表 updated_gmt_at(字典表/contacts 无更新列不列)

验证:tsc 0 err;jvs-dw 零影响(manifest 未动 + ?? 短路 + dry-run Cursor 解析行为不变);
friday 实跑 sync_logs 首次产出 cursor_after={12 表 × run_start ISO}。
心智模型统一:所有宿主存量跑即水位创世写入,差别只剩水位消费方。

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
parent 8b556c4a
......@@ -27,6 +27,25 @@ identity_namespace_field: brand_id # patients.source_unit ← 源品牌
amount_unit: yuan
timezone: Asia/Shanghai
# ── 增量水位声明(2026-07 与 jvs-dw 统一;file 源放 manifest 顶层)──
# file 装载不消费水位(每轮全量装载,幂等去重兜底);每轮跑完与 jvs-dw 走同一段 finally
# 记账 cursor_after = run_start(sync_logs)。用途:delta 导出 WHERE 模板 / push 回放起点。
# 只列有 updated_gmt_at 的事件表;字典表(std_*)与 customer_contacts(无更新列)恒全量。
incremental:
per_query:
customer_basic_info: { cursor_column: updated_gmt_at }
appointment_base: { cursor_column: updated_gmt_at }
med_emr_info: { cursor_column: updated_gmt_at }
med_check: { cursor_column: updated_gmt_at }
patient_settlement: { cursor_column: updated_gmt_at }
patient_settlement_refund: { cursor_column: updated_gmt_at }
patient_settlement_spec_refund: { cursor_column: updated_gmt_at }
settlement_modes: { cursor_column: updated_gmt_at }
customer_treat_plan: { cursor_column: updated_gmt_at }
customer_treat_plan_item: { cursor_column: updated_gmt_at }
customer_referee_circle: { cursor_column: updated_gmt_at }
customer_consult: { cursor_column: updated_gmt_at }
tables:
- { table: customer_basic_info, file: customer_basic_info.csv } # 患者主档(原始)
- { table: customer_contacts, file: customer_contacts.csv } # 联系方式 1:N(原始)
......
......@@ -639,7 +639,12 @@ export class ColdImportService {
// 避免推进到 run_start 后把 DW 迟到的旧 updated_date 行永久埋在游标下方(静默漏拉)。
let incrementalNoData = false;
let incrementalBaselineCursor: Record<string, string> = {};
const perQueryCfg = manifest.sql_source?.incremental?.per_query;
// 水位声明与源类型解耦(2026-07 统一):sql_source 内优先(jvs-dw 现状不动),
// file 源用 manifest 顶层声明 —— 之后读写 cursor 走完全同一段逻辑(finally 记账
// cursor_after = run_start)。file 装载路径不消费水位(loadAllTables 文件分支忽略
// incremental 参数,全量装载靠幂等去重),水位仅供 delta 导出 / push 回放起点。
const perQueryCfg =
manifest.sql_source?.incremental?.per_query ?? manifest.incremental?.per_query;
if (perQueryCfg) {
// 默认读 cursor;options.incremental === false 时 (强制 full) 不读
const ignoreCursor = options.incremental === false;
......@@ -681,7 +686,7 @@ export class ColdImportService {
} else if (options.incremental) {
// 老调用方明确要 incremental 但 manifest 没配 cursor → fast-fail
throw new Error(
`incremental 模式需要 manifest.sql_source.incremental.per_query;请补 cursor 配置`,
`incremental 模式需要 incremental.per_query 声明(sql_source 内或 manifest 顶层);请补 cursor 配置`,
);
}
......
......@@ -47,6 +47,22 @@ export type AssemblerRef = z.infer<typeof AssemblerRefSchema>;
*
* 密码:从环境变量读取(`password_env`),不允许 plaintext 进 yaml。
*/
/// 增量水位声明:表名 → cursor 列。**与源类型解耦**(2026-07 统一):
/// - sql_source 内声明(jvs-dw 现状,位置不动):拉模式,读水位注 WHERE + 写水位
/// - manifest 顶层声明(file 源用):file 装载不消费水位(全量装载,幂等去重兜底),
/// 只做**记账**——同一段 finally 写 cursor_after = run_start,水位供 delta 导出
/// WHERE 模板 / push 回放起点使用。两处都写时 sql_source 内优先。
export const IncrementalDeclSchema = z.object({
/// per query/table 配置;表名 → cursor_column
/// 例:fact_emr_treatment_out: updated_date / med_emr_info: updated_gmt_at
per_query: z.record(
z.string(),
z.object({
cursor_column: z.string().min(1),
}),
),
});
export const ClickHouseSourceSchema = z.object({
kind: z.literal('clickhouse'),
connection: z.object({
......@@ -65,20 +81,9 @@ export const ClickHouseSourceSchema = z.object({
default_limit: z.number().int().positive().default(100000).optional(),
/// W4 末:DW 增量配置(per table cursor column)
/// 跑增量模式时 PAC 会读 sync_logs 上次 cursor_after,注入 WHERE cursor_column > '...'
/// 跑完写新 cursor = max(cursor_column) 到 sync_logs
/// 跑完写新 cursor = run_start 到 sync_logs
/// 首次跑(无 cursor)= 全量,后续 = 增量
incremental: z
.object({
/// per query 配置;表名(query key)→ cursor_column
/// 例:fact_emr_treatment_out: updated_date / fact_client_out: last_visit_time
per_query: z.record(
z.string(),
z.object({
cursor_column: z.string().min(1),
}),
),
})
.optional(),
incremental: IncrementalDeclSchema.optional(),
/// ⭐ Cohort 分批配置(PR2/PR4)— 宿主无关性的关键:
/// patient 列表来源 + 主键列 + 租户区分列全部声明在这,代码不硬编码任何表名/列名。
/// 不配 cohort → 不分批(single-shot,文件源 / 小数据集)。
......@@ -179,6 +184,11 @@ export const ColdImportManifestSchema = z
/// ClickHouse SQL 数据源(跟 tables 二选一)
sql_source: ClickHouseSourceSchema.optional(),
/// 顶层增量水位声明 —— **file 源专用**(sql_source 宿主继续在 sql_source.incremental 声明,
/// 两处都写时 sql_source 内优先)。file 装载不消费水位,仅记账 cursor_after = run_start
/// (同 sql_source 宿主同一段 finally,逻辑一致);详见 IncrementalDeclSchema 注释。
incremental: IncrementalDeclSchema.optional(),
/// 诊所名字典来源(可选)—— 声明哪张表的哪两列 = 诊所 id → 名字。
/// `refresh-clinic-names` CLI 据此 SELECT DISTINCT 出 id→名,写 host.clinicNames;
/// /auth/session 再合并进 dictionary.clinics 下发,前端显示名字不依赖登录传。
......
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