Nango Records 数据迁移实战:从 _nango_sync_data_records 到 nango_records 的断点续传迁移与一致性校验

Nango Records 数据迁移实战:从 _nango_sync_data_records 到 nango_records 的断点续传迁移与一致性校验 Nango Records 数据迁移实战从 _nango_sync_data_records 到 nango_records 的断点续传迁移与一致性校验【免费下载链接】nangoBuild product integrations with AI.项目地址: https://gitcode.com/GitHub_Trending/na/nango本指南基于 Nango 仓库中 scripts/one-off/records-migration 目录下的运维工具完整讲解如何把历史同步数据记录从旧的nango._nango_sync_data_records表迁移到新的nango_records.records表。读完本文你将掌握该迁移工具的配置方法、批量迁移与断点续传原理、以及如何用配套的verify.ts脚本对迁移结果做逐行一致性校验可直接复用到你自己的数据搬迁场景。一、脚本定位与适用场景records-migration是 Nango 仓库中scripts/one-off下的一个独立运维脚本用于跨库/跨 Schema 迁移同步数据记录sync data records。其目录结构如下README.md运行说明本文主体migrate.ts数据迁移主脚本verify.ts迁移结果一致性校验脚本package.json依赖与 npm scriptstsconfig.jsonTypeScript 编译配置它解决的典型问题是当 Nango 内部把同步记录从旧表nango._nango_sync_data_records位于nangoschema迁移到新表nango_records.records位于nango_recordsschema时需要在两个 PostgreSQL 数据库/实例之间搬运海量行数据同时保证不丢数据、可中断恢复、且迁移后数据与原库一致。二、快速开始README 中的运行方式README 给出的运行步骤非常精简只有三步npm install npm run migrate从源码实现看完整的运行前提是配置数据库连接通过环境变量SOURCE_DB_URL与TARGET_DB_URL分别指定源库与目标库的连接串。README 中写的是UpdateSOURCE_DB_URLandTARGET_DB_URLinto./migrate.ts而实际代码migrate.ts是从process.env读取这两个变量的因此更推荐也更符合真实行为使用环境变量注入连接串export SOURCE_DB_URLpostgresql://user:passsource-host:5432/nango export TARGET_DB_URLpostgresql://user:passtarget-host:5432/nango npm run migrate安装依赖npm install只需安装一个核心依赖 knex^3.1.0见 package.json。执行迁移npm run migrate实际执行的是npm run build node ./dist/migrate.js即先用tsc把 TypeScript 编译到dist/目录再以 Node 运行产物。同理npm run verify对应npm run build node ./dist/verify.js。脚本按 tsconfig.json 编译module为es2022、moduleResolution为Bundler、rootDir为.、输出到dist且 package.json 声明了type: module因此源码中使用import.meta.url等 ESM 特性。三、连接配置与连接池参数详解两个脚本的数据库连接配置完全一致migrate.ts 与 verify.ts以 Knex pg客户端连接 PostgreSQL配置项取值说明clientpgPostgreSQL 驱动connection.connectionStringSOURCE_DB_URL/TARGET_DB_URL数据库连接串从环境变量读取connection.sslno-verify启用 SSL 但跳过证书校验适配云数据库场景connection.statement_timeout60000毫秒单条 SQL 语句最长执行 60 秒防止长查询挂死pool.min/pool.max1/10连接池最小 1、最大 10 个连接两点差异需要注意migrate.ts中两个环境变量默认值为空字符串不设置会导致连接失败verify.ts中则提供了本地开发默认值postgresql://nango:nangolocalhost:5432/nango方便在本地自建 PostgreSQL 时直接运行校验。四、迁移主脚本 migrate.ts 的实现原理migrate.ts是整个搬迁的核心其设计要点可以拆解为以下五个机制。4.1 源表字段映射readRecordsmigrate.ts从nango._nango_sync_data_records表按批读取并在 SQL 层完成字段重命名与类型转换id, external_id, json, data_hash, nango_connection_id as connection_id, model, sync_id, sync_job_id, to_json(created_at) as created_at, to_json(updated_at) as updated_at, to_json(external_deleted_at) as deleted_at, created_at as created_at_raw关键映射关系源表nango_connection_id→ 目标表connection_id源表external_deleted_at软删除时间戳→ 目标表deleted_atcreated_at/updated_at用to_json()转换为 JSON 文本同时保留原始created_at为created_at_raw供排序分页使用。4.2 键集分页Keyset Pagination迁移按批大小 1000 行BATCH_SIZE读取采用(created_at, id)组合键做键集分页而不是传统的OFFSETbuilder.where(sourceKnex.raw((created_at, id) (?, ?), [checkpoint.lastCreatedAt, checkpoint.lastId]));查询按created_at_raw ASC, id ASC排序并LIMIT 1000。这种分页方式的优势是即便迁移过程中源表有新数据写入也不会像OFFSET那样出现跳行或重复天然适合持续跑批的大表迁移。4.3 断点续传Checkpoint脚本会在脚本同目录生成migrate_checkpoint.json文件migrate.ts每次成功写入一批后记录{ lastCreatedAt: ..., lastId: ... }启动时通过getCheckpoint()读取该文件若文件不存在ENOENT则从{ lastCreatedAt: null, lastId: null }全量开始migrate.ts每批迁移完成后立即saveCheckpoint()落盘若中途进程崩溃或网络中断重新运行npm run migrate即可从断点继续已迁移的数据不会重复处理。4.4 写入采用 Upsert 合并writeRecordsmigrate.ts向目标表nango_records.records执行批量插入冲突时按业务唯一键合并targetKnex .insert(records) .into(nango_records.records) .onConflict([connection_id, model, external_id]) .merge() .returning([id, created_at]);即以(connection_id, model, external_id)作为业务唯一键若目标表已存在同键记录则更新覆盖并返回实际写入的id与created_at。这保证了迁移的幂等性——重复执行不会产生重复行。4.5 读写流水线与空批次处理迁移主循环migrate.ts做了两个优化写读并行Promise.all([writeRecords(toInsert), readRecords(checkpoint)])在写入当前批次的同时预读下一批次用连接池双连接摊薄网络延迟空批次休眠当records.length 0时打印No rows to migrate. Sleeping...并休眠 2 秒后重试而不是直接退出适合源表仍在持续写入、需要追尾的场景。每批完成后控制台输出类似1000 rows migrated in 350ms. lastCreatedAt: 2024-01-01T00:00:00.000Z.脚本结束时会销毁两个 Knex 连接池sourceKnex.destroy()/targetKnex.destroy()并打印总耗时migrate.ts。若中途发生异常会打印Error occurred during data migration并保留 checkpoint 文件以便续跑。五、一致性校验脚本 verify.ts迁移完成后verify.ts 负责把目标表数据逐批拉回来与源表比对找出不一致的记录。5.1 分区扫描校验脚本假定目标表按某种规则分为 256 个分区表名为nango_records.records_p0到nango_records.records_p255外层for循环遍历每个分区verify.ts。内层同样以 1000 行为一批但分页键换成了(updated_at, id)并按updated_at_raw ASC, id ASC排序checkpoint 存于verify_checkpoint.json。5.2 差异比对逻辑每批读取后getDiffsverify.ts用 PostgreSQL 的json_populate_recordset把内存中的记录转成临时记录集再LEFT JOIN回源表nango._nango_sync_data_records逐字段比对行缺失r.id IS NULL目标存在但源表找不到对应(external_id, model, connection_id)字段不一致id、data_hash、sync_id任一不同时间戳漂移created_at/updated_at与源表相差超过30 秒ABS(EXTRACT(EPOCH FROM ...)) 30判定为不一致容忍毫秒级或时区换算差异删除状态不一致目标deleted_at与源表external_deleted_at一个为空一个非空或时间差超过 30 秒。比对时先剔除updated_at_raw与json字段JSON 全文不参与比对以data_hash为准避免噪音。每发现差异会打印前 5 条样例并累计diffsCountFound 3 diffs (total: 3)): [ ...前5条样例... ] 1234 rows verified (parition 0) in 220ms. lastUpdatedAt: ...5.3 校验结论全部 256 个分区扫完后输出Verification completed. Total diffs found: 0当Total diffs found为 0 时即可认为迁移后的目标表与源表在业务语义上完全一致可以放心切换读写路径。六、源表结构佐证数据库迁移文件Nango 仓库的数据库迁移文件证实了源表_nango_sync_data_records的结构与迁移脚本的字段假设完全吻合。创建该表的迁移为 packages/database/lib/migrations/20230519080408_add_sync_data.cjs核心列包括iduuid非空external_idstring非空jsonjsonb可为空data_hashstring非空nango_connection_idinteger外键指向_nango_connections.idON DELETE CASCADEmodelstring非空created_at/updated_attimestamp唯一约束unique([nango_connection_id, external_id])此外该迁移还创建了update__nango_sync_data_records_updated_at()触发器函数仅在data_hash变化时才自动刷新updated_at这也是校验脚本以data_hash作为内容指纹、以updated_at作为变更分页键的原因。软删除列external_deleted_at则由 packages/database/lib/migrations/20230810110751_data_records_deleted.cjs 等后续迁移加入对应迁移脚本中to_json(external_deleted_at) as deleted_at的映射。七、实战注意事项与使用边界基于源码实现使用该工具时有几点需要特别留意README 与代码的差异README 说把连接串写进 migrate.ts但代码实际读取的是环境变量。以代码为准用环境变量注入连接串即可无需改文件。目标表需预先就绪nango_records.records及校验用的分区表records_p0..p255需要已存在且具有(connection_id, model, external_id)的唯一约束否则onConflict无法生效、校验脚本也会报错。断点续传靠本地文件migrate_checkpoint.json与verify_checkpoint.json存放在脚本目录下换机器或清理文件后将从零开始。如果源表在迁移期间仍有写入且你希望全量重跑请先删除 checkpoint 文件再执行。校验脚本的分区假设verify.ts硬编码了 256 个分区的遍历若实际部署的分区数不同需要相应调整外层循环。迁移是追尾式而非终止式migrate 在读到空批次时休眠 2 秒后继续适合源表持续有增量写入的场景若要停表迁移可在源表写入停止后观察输出不再产生新行时手动终止。八、总结records-migration是一套小而完整的可断点续传 幂等 upsert 独立一致性校验的数据搬迁方案migrate.ts用键集分页与 checkpoint 保证海量数据的可靠搬运verify.ts用json_populate_recordsetLEFT JOIN做逐行差异比对配合data_hash指纹与 30 秒时间容差过滤噪音。其设计模式环境变量注入连接、批量流水线、断点文件、业务唯一键合并、独立校验脚本可以平移到其他任何需要跨库迁移海量业务数据的运维任务中。【免费下载链接】nangoBuild product integrations with AI.项目地址: https://gitcode.com/GitHub_Trending/na/nango创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考