最近在一个业务系统的数据接入项目里,我被一个看似简单的问题折磨了好几天:源端的订单表和用户表几乎每十分钟就会有一批新数据进来,高峰时段一小时能冒出几千条新记录。我用 n8n 搭了一个同步工作流,第一版图省事直接做全量拉取,结果跑了两天就扛不住了——数据源 API 开始限流报警,目标库频繁出现锁等待,整个流程又慢又脆。后来我花了整整一个周末,把增量同步的完整逻辑重新设计了一遍,才终于把问题理顺。
这篇文章我不打算讲那些“官方文档里都写了”的基础用法,而是想把我踩过的坑、最终沉淀下来的方案,以及 n8n 里设计增量同步工作流时真正需要注意的细节,完整地写出来。如果你手里也有一套数据源动不动就更新、靠全量同步凑合了挺久的工作流,这篇文章应该能帮你少走不少弯路。
1. 先搞清楚:你面临的增量同步到底是什么问题
1.1 全量同步为什么总是撑不住
全量同步的思路其实很朴素:每次定时任务触发,就把源端整张表、整个列表全部拉回来,然后整体覆盖到目标端。数据量小的时候,这种做法确实省心,一行SELECT * FROM orders就能搞定,n8n 里搭一个定时触发器再连一个写入节点,十几分钟就能上线。
但这套逻辑一旦碰上频繁更新的数据源,问题会像滚雪球一样冒出来。首先是数据量本身在膨胀,每天都有新增长,全量拉取的数据包会越来越大,n8n 的执行节点需要逐步把整批数据装载到内存里,数据行数一多,执行的耗时和内存占用都会直线上升。其次是源端压力,你会发现数据源 API 并不是无限量供应,高频的全量请求很容易把接口配额打穿,对方直接给你返回 429。最后是目标端的写入压力,每次全量同步都相当于把全表删掉再重建一遍,锁等待、主键冲突、性能抖动全都跟着来了。
我当时就把全量同步改成每半小时跑一次,结果源端数据库的慢查询日志里几乎全是我的同步语句,运维同事直接找上门问我到底在干什么。全量同步不一定是错误的,但它天然不适合“数据量持续增长 + 更新频率高”这个组合。
1.2 增量同步的本质不是“只拉最近十分钟”
很多人一听增量同步,第一反应是“那我加个时间条件,只查最近十分钟的数据不就行了”。真这么干,大概率会掉进另一个坑。
增量同步的本质,是确定一个同步水位线(watermark)。你要维护一个“上一次已经同步到哪里了”的边界,下一次运行时,只处理边界之后出现或变化的数据。这个边界可以是时间戳,可以是自增 ID,也可以是某个事务日志里的位置。水位线必须稳定、可追溯,并且能抵抗部分失败——如果这次同步跑到一半挂了,下一次重跑时不能漏数据,也不能因为重复推进一步而丢数据。
我们在 n8n 里做增量同步,本质上就是要回答三个问题:
- 水位线存在哪里?
- 下一次运行时,如何把水位线精确地应用到查询或 API 请求里?
- 一批数据处理完以后,水位线应该如何安全推进?
这三个问题能答好,增量同步的核心骨架就立住了。后面的策略选择、节点编排、异常处理,全都是围绕这三个问题展开的。
2. 增量同步的四种主流思路与选型逻辑
2.1 基于修改时间戳的方案:最直观也最容易踩坑
这是最常见的做法,前提是源端表里有一个“最后修改时间”字段,比如updated_at、modify_time。逻辑很简单:同步时查询updated_at > 上次水位线的数据,处理完之后把当前时间或者这批数据里的最大updated_at更新为新的水位线。
听起来简单,实际项目里有几个坑要提前想清楚。
第一个坑是时间精度。很多数据库的时间字段只精确到秒,如果一个事务里同时更新了一百条记录,它们的updated_at完全可能是同一个值。如果你查询条件用的是严格大于>,这批记录刚好和上次水位线撞在同一秒,就会被漏掉。更稳妥的方案是用>=,并且把主键作为第二排序条件,保证同秒记录也能被完整捞出来。
SELECT id, title, content, updated_at FROM articles WHERE updated_at >= :lastSyncAt ORDER BY updated_at ASC, id ASC LIMIT 500;第二个坑是时间归属问题。源端数据库的时区设置、应用写入时的时间戳转换,都可能导致“看似合理的时间条件”失效。我自己就被坑过一次:源端库用的Asia/Shanghai,业务写入时却把 UTC 时间直接塞进了字段,结果增量同步经常漏掉下午更新过的数据。处理这类问题的通用原则是:同步链路里所有时间的生产、比较、存储,统一用 UTC,字段名里明确标注,不要用服务器本地时间。
第三个坑是索引。updated_at字段如果没有索引,增量查询依然会演变成全表扫描。特别是当你只拉取几万条里的几百条时,有无索引的性能差距是数量级的。建议在源端给(updated_at, id)建一个组合索引,既能过滤时间范围,又能配合排序。
2.2 基于自增 ID / 最大主键的方案:只适合追加型数据
如果你的数据源是一个只增不改的结构,比如操作日志、点击流、订单创建记录,那么用“最大自增 ID”作为水位线会非常舒服。实现更简单:记录上次同步的最大 ID,下次查询时WHERE id > :lastMaxId,配合分页把数据拉完。
这个方案最大的优势是精准,自增 ID 天然单调递增,不会像时间戳那样出现边界模糊。但它的短板同样明显:完全无法感知历史记录的修改和删除。源端如果存在“插入后又被回改”的场景,比如用户先下了一单然后又取消,订单状态从pending改成了closed,这个变更不会产生新的自增 ID,增量同步就会漏掉。
所以我通常在选型时会画一条分界线:如果源端表只做插入、不做更新,或者更新行为不影响同步目标,就可以用自增 ID;只要存在任何更新历史记录的业务场景,就老老实实回退到时间戳方案,或者直接把两套方案结合着用。
2.3 基于事务日志或 CDC 的方案:真高频数据源的终极答案
当同步频率要求极高,比如秒级延迟、分钟级延迟,并且源端体量大到不能容忍任何全表或者宽时间范围的扫描时,基于时间戳的方案也开始不够看了。这时候业界的主流做法是引入 CDC(Change Data Capture),直接读取数据库的 binlog 或者 WAL 日志,把每一条增删改操作都解析成事件流。
n8n 本身不是 CDC 工具,但完全可以作为 CDC 事件流的消费端。典型链路是:源数据库的日志被 Debezium 这样的工具解析后,推送到消息队列,n8n 提供一个 Webhook 端点,每当有数据变更事件到达,就触发工作流,把变更应用到目标库。
这套方案很强大,但它也确实重。你要额外维护一个 CDC 解析服务、一个消息队列,还要考虑 schema 变更、事件顺序、重复投递等问题。我个人把它定位为“重武器”,只有在数据频率和规模真的到了量级,并且团队有能力维护额外基础设施时,才会考虑它。绝大多数中小型项目,做好时间戳增量配合幂等写入,已经能覆盖 90% 以上的业务场景。
2.4 基于 API 自身增量能力的方案:千万别自造轮子
有些 SaaS API 或者内部服务的接口,本身就提供了“只返回某个时间点之后变更的数据”的能力。比如像 Shopify 的updated_at_min参数、各类 CMS 的modified_since头,或者某些平台直接给出基于游标的分页接口,游标本身就是水位的体现。
遇到这类接口,优先直接使用原生能力,别自己绕路。把 API 参数里的时间范围和游标直接映射到 n8n 工作流里,会省掉大量过滤和数据比对工作。不过也要注意几个细节:API 是否有分页上限,比如单页只能返回 250 条;时间字段是否允许精确到毫秒;API 返回的记录排序是否稳定。我见过一个同事的同步工作流经常跳数据,排查了半天,发现是 API 的分页排序没有主键兜底,下一页和上一页之间偶发重叠。这种情况直接在请求参数里加上排序字段,就能解决。
3. n8n 里增量同步的关键设计:状态记忆
3.1 用轻量数据库表记录同步水位线
不少 n8n 新手会把水位线存在工作流变量里,或者干脆硬写在节点配置中,这是非常容易出事的做法。n8n 的工作流在每次执行时,节点之间的数据只会存在于当次执行的上下文里,下一次执行时上一轮的变量不会天然保留。全局变量功能在部分场景下可用,但生产级的同步任务,我更推荐用一个独立的表来维护水位线。
这个表不需要复杂,三五个字段足够:
CREATE TABLE sync_state ( id SERIAL PRIMARY KEY, source_name VARCHAR(255) NOT NULL UNIQUE, last_cursor_value TIMESTAMP NOT NULL, last_max_id BIGINT, updated_at TIMESTAMP DEFAULT now() );用这个表的好处非常直接:水位线是持久化的,即使 n8n 容器重启、工作流被重新部署、某个流程跑挂了,水位线也不会丢。而且在错误诊断时,你可以直接查询这张表,看到每个数据源当前推进到了什么位置,很多“数据到底同步到哪了”的争论一眼就能定位。
3.2 把水位线精确传给查询参数
在 n8n 工作流里,读取水位线的位置通常放在执行链路的头部。触发节点跑起来之后,先连一个 Postgres 节点,执行类似下面的查询:
SELECT current_timestamp AS default_cursor FROM sync_state WHERE source_name = 'articles';查询结果会变成后续节点的输入,你会在后面的代码节点或 HTTP Request 节点里通过{{ $json.last_cursor_value }}引用这个值。这里要注意一个细节:如果sync_state里暂时还没有对应数据源的记录,也就是首次运行时,查询会返回空结果,后续节点会直接报错。解决办法是给这个查询用COALESCE或者UNION兜一个默认值,比如第一次运行时允许回溯到七天前,保证第一次同步也能有数据。
SELECT COALESCE( (SELECT last_cursor_value FROM sync_state WHERE source_name = 'articles'), (CURRENT_TIMESTAMP - INTERVAL '7 days') ) AS cursor_value;这样不管是有状态还是无状态,查询结果两边的字段名都是一致的,后续节点不用为“首次运行”和“日常运行”分别写两套逻辑。
3.3 数据分批与循环:别让一次请求扛下所有
增量同步的设计里,水位的读取只是开始,真正复杂的其实是“一次拉取多少数据”和“怎么把多批数据拼起来”。如果你直接查询全部增量数据,比如一次性SELECT * FROM articles WHERE updated_at >= '2024-01-01',数据量照样可能达到几万甚至几十万行,内存压力又会回来。
更合理的做法是引入批次拉取。比如每批 500 行,用LIMIT 500 OFFSET n或者通过排序和上一批最大 ID 来翻页。在 n8n 里常见的编排方式是“循环节点 + 合并节点”:每次循环处理一批数据,批处理完成后把结果追加到同一个数组里,直到当前批次返回的行数不足一批,就停止循环。
这里有一个我反复强调的点:确认“还有没有下一批”的判断条件,不能只看是否等于批次大小。比如设定每批 500 行,当返回结果恰好等于 500 时,你以为还有更多数据,其实可能刚好就是最后一批。稳妥的做法是让查询语句多取一行,比如LIMIT 501,如果返回结果是 501 行,说明确实还有下一批;如果只有 500 行或更少,那就说明数据拉完了,这一批里的最后一行在下一轮才对,从而处理实际情况。
这么设计的好处是,即使中间某次执行失败,下一个执行周期的水位线没有推进,重跑时用>=只会造成少量重复数据,不会造成缺失。配合后续讲到的幂等写入,重复数据本身也不是致命问题。
3.4 一个最小可用的同步链路长什么样
把上面的思路归纳一下,一个最小可用的 n8n 增量同步工作流,节点编排大致是这样:
- Schedule Trigger:每几分钟或每小时触发一次。
- Postgres 节点:读取
sync_state里的水位线。 - Code 节点:组装请求参数或查询 SQL,设置批次大小和排序规则。
- HTTP Request 节点 或 Postgres 查询节点:拉取一批增量数据。
- Code 节点 或 Condition 节点:判断这批数据是否还有下一页,有则继续循环。
- 数据清洗节点:统一字段名、处理缺失值、转换时间格式。
- 目标写入节点:用 upsert 方式把数据写进目标库。
- 状态推进节点:整批数据全部成功后,更新
sync_state里的水位线。
这套链路的顺序是有讲究的。很多人喜欢在循环最开始就推进水位线,发现性能问题或者逻辑漏洞时再回头改,往往要重写一大片。核心原则是:水位线只升不降,并且永远在处理完本批数据并成功写入之后才推进。如果写入失败了,水位线保持在原位,下一次重跑时会从更早的位置从头处理,配合幂等写入,相当于自动把失败补齐了。
4. 在 n8n 中逐步实现一个时间戳增量同步
4.1 第一步:搭建触发器与初始化
我以 PostgreSQL 数据源为例子,详细拆解一下每一步怎么在 n8n 中落地。
首先拖一个 Schedule Trigger 节点,配置成Run Workflow Every10 分钟,或者更贴合业务的高峰时段可以改成 1 分钟。注意,n8n 在生产环境跑定时任务时,触发节点本身并不保证任务执行的耗时不会影响下一次触发,如果你同步的数据量大、执行时间可能超过触发间隔,就需要在 Schedule Trigger 里把时区配置准确,并适当增大间隔,避免多个实例重叠执行。
首次运行前,先把sync_state表建好,并插入一条初始记录,水位线可以为空,查询逻辑用COALESCE兜底:
INSERT INTO sync_state (source_name, last_cursor_value) VALUES ('articles', NULL) ON CONFLICT (source_name) DO NOTHING;4.2 第二步:读取上次游标并组装查询
触发器后面接一个 Postgres 节点,查询水位线。SQL 里使用COALESCE让你在表里没有记录时也能拿到一个默认的同步起点。
SELECT COALESCE( max(last_cursor_value), current_timestamp - interval '3 days' ) AS cursor_value FROM sync_state WHERE source_name = 'articles';查询结果出来以后,接一个 Code 节点,把上游返回的游标值塞进一个统一的对象里,方便后续节点引用。比如:
const cursorValue = $json.cursor_value; return [{ cursorValue: cursorValue, batchSize: 500 }];这里我特意把上游字段名在 Code 节点里统一成cursorValue,目的有两个:第一,SQL 查询的结果字段可能因为不同数据库驱动而大小写不一致,代码节点里统一命名能避免后续节点到处踩字段名大小写的坑;第二,后续如果要切换不同的数据源,只需要改这一个节点,后面所有节点都不用动。
4.3 第三步:字段筛选与分页循环
拿到游标后,就可以去业务表里拉增量数据了。这一步在真实项目里我几乎都会在 SQL 里直接做字段筛选,不要偷懒用SELECT *,一方面是减少网络传输量,另一方面是避免把一些二进制大字段或者敏感字段带出来。
SELECT id, title, status, updated_at FROM articles WHERE updated_at >= '{{ $json.cursorValue }}' ORDER BY updated_at ASC, id ASC LIMIT {{ $json.batchSize }};分页循环我建议使用 n8n 的 Loop Over Items 节点,把上一批查询到的最后一条记录的updated_at和id作为下一批查询的起点。具体做法是在循环内部用一个 Postgres 节点执行查询,查询条件动态变成:
WHERE (updated_at > :lastUpdatedAt) OR (updated_at = :lastUpdatedAt AND id > :lastId)这种写法替代OFFSET翻页有两个好处。第一,不会因为数据量大导致OFFSET越来越大、查询越来越慢;第二,如果循环过程中有新的数据插入,也不会出现翻页时跳过记录的情况。这套翻页方式通常叫键集分页(keyset pagination),是增量同步领域里非常实用的技巧。
循环的终止条件就看返回行数是否小于批次大小,等于批次大小时继续循环下一轮,直到行数小于批次大小才跳出。
4.4 第四步:数据清洗与目标写入
循环里把数据一批批拉出来后,接到一个 Code 节点,对字段做统一清洗。这里的常见问题包括:时间字段从字符串转成标准格式、空字符串转为 NULL、源端是枚举字段而目标端需要映射成不同的值。清洗逻辑写得越多,目标端的写入就越省心。
清洗完之后,需要把多批数据攒到一起再写目标端,还是每批单独写?我的经验是:写入频率不能太频繁,也不能一次性堆太大。每批 500 行,边拉边写,既不会每一条都触发一次网络往返,又不会让内存暴涨。清洗完一批,就直接用一个 Postgres 节点执行 upsert:
INSERT INTO articles_sync ( id, title, status, updated_at ) VALUES ( $1, $2, $3, $4 ) ON CONFLICT (id) DO UPDATE SET title = EXCLUDED.title, status = EXCLUDED.status, updated_at = EXCLUDED.updated_at;使用 upsert 的核心原因是增量同步过程中不可避免存在重复数据和乱序数据。数据源可能因为网络重试、分页重叠,把同一条记录传送了两次;也可能因为事务提交顺序不同,导致目标端先收到早先版本,后收到更新版本。用ON CONFLICT (id) DO UPDATE,按主键去重并且覆盖更新,就能把这些脏数据都吸收掉。
目标端写入后,我还习惯接一个 Count 节点或者 Code 节点记录一下本批写入的行数,方便后面做监控和排查。零行同步和大量异常同步,日志里必须要能区分。
4.5 第五步:成功后再推进水位线
整批数据全部写完,最后一步才是更新sync_state表:
UPDATE sync_state SET last_cursor_value = :lastUpdatedAt, updated_at = now() WHERE source_name = 'articles';推进的游标值用这一批数据里的最大updated_at,而不是“当前系统时间”。直接用系统时间作为水位线的坑在于,如果数据源时钟和本地时钟存在偏差,或者业务数据里某些记录的updated_at是手动填写的未来时间,那么本批数据里可能还有不少记录的时间大于当前时间,却被错误排除在下一批范围之外。用数据中的最大时间戳作为游标,能够保证这条边界始终和数据本身的分布对齐。
5. 实际运行中必须处理的复杂情况
5.1 删除、软删除与墓碑记录
基于时间戳的增量同步,有一个天生的盲区:它看不到删除。如果源端直接物理删除了某一行,目标端只能毫不知情地保留一条已经不存在的数据。
遇到这种情况,我会先看源端有没有软删除设计。如果业务表里有一个deleted_at字段,那情况就简单了,同步时把deleted_at也算进updated_at的逻辑里,凡是deleted_at不为空的行,增量拉过来后,目标端除了更新普通字段外,还要根据deleted_at做删除或者标记。
SELECT id, title, status, deleted_at, updated_at FROM articles WHERE updated_at >= :lastSyncAt OR deleted_at >= :lastSyncAt ORDER BY updated_at ASC, id ASC;如果源端本身完全不保留删除痕迹,我会退而求其次,保留一个“对账周期”。比如每周做一次全量比对,把两边主键集合做差,把已经不在源端的记录在目标端做逻辑删除。这虽然听起来很笨,但在没有 CDC 和可靠事件流的情况下,已经是最稳妥的兜底方案。
5.2 重复数据与乱序写入
前面讲了 upsert 能吸收重复数据,但对乱序写入,还有一个更隐蔽的问题:如果源端更新了多条记录,其中一条早先版本的写入晚于另一个新版本,目标端可能被旧数据覆盖。
既然无法保证源端消息顺序,我的防线是加一个source_updated_at字段,专门用来记录源端这次修改的时间。目标端写入时,不只是简单覆盖,而是先判断新来的source_updated_at是否比目标端已有记录的新,只有新的才覆盖:
INSERT INTO articles_sync (...) VALUES (...) ON CONFLICT (id) DO UPDATE SET title = CASE WHEN EXCLUDED.source_updated_at > articles_sync.source_updated_at THEN EXCLUDED.title ELSE articles_sync.title END, source_updated_at = GREATEST(articles_sync.source_updated_at, EXCLUDED.source_updated_at);这种做法本质上是在目标端做了一次基于时间戳的“乐观锁”,虽然 SQL 看起来复杂一点,但它能彻底杜绝乱序写入导致的数据回退问题。这个坑我踩过不止一次,特别是一次迁移后连续好几天的数据看起来没问题,实际上一些核心资料已经被旧版本覆盖,等发现时已经很难追溯是哪一次写入造成的了。
5.3 数据源偶尔返回超大分页
有些数据源虽然支持分页,但单页最大记录数可能设为 1000 或者 2000。当增量窗口内有大量数据更新时,你会遇到拉取一批数据特别多、循环体执行时间特别长的情况。这时候 n8n 节点的内存控制和执行超时就成了潜在隐患。
我的处理方式是设置一个最大批次保护。查询 SQL 里,如果分页没有明确限制,就在代码里强行限制每次最多处理 500 行,即使数据源允许每页 2000 行,也不要真的去拿那么多。更大的批次虽然减少网络往返,但会让 n8n 临时内存里堆积的对象暴增,一旦执行节点所在的容器内存吃紧,整个工作流都可能被 OOM 干掉,反而比多循环几次更危险。
另外,n8n 里如果使用循环收集结果,记得在循环结束后用 Split Into Batches 节点把结果集拆成合适大小,再分批写入目标端。永远不要在同一个内存数组里堆几十万行记录。
5.4 失败追踪与重试
增量同步工作流跑久了,必然会遇到失败。失败并不可怕,可怕的是失败后没有任何线索,只能靠人工翻日志。
我在 n8n 里的做法是单独建一个同步日志表:
CREATE TABLE sync_runs ( id SERIAL PRIMARY KEY, source_name VARCHAR(255) NOT NULL, started_at TIMESTAMP NOT NULL, finished_at TIMESTAMP, status VARCHAR(20) NOT NULL, rows_processed INT, error_message TEXT );工作流开始时插入一条status='running'的记录,成功结束时更新为success并把处理行数写上,异常分支里则写入failed和错误明细。n8n 的 Error Trigger 节点可以在工作流失败时捕获错误信息,把错误对象透传给这个日志表。
除此之外,还要给关键数据源配上失败通知,我习惯用一个分支节点,让失败记录通过 Telegram 节点或者企业微信机器人发到群里。这一步虽然不起眼,但当你半夜被数据同步失败的电话吵醒时,一条带有错误上下文的群消息,能帮你把排查时间从一小时压缩到五分钟。
6. 我的几点实操心得与排错经验
6.1 常见问题速查表
| 问题现象 | 可能原因 | 解决方法 |
|---|---|---|
| 增量同步漏掉同秒更新的记录 | 查询边界用了>,且时间精度到秒 | 改为>=,搭配主键排序和键集翻页 |
| 目标端数据被旧版本覆盖 | 乱序写入,新版本先到、旧版本后到 | upsert 里对比source_updated_at,只允许更新的值覆盖 |
| 同步任务越跑越慢 | updated_at字段没有索引,翻页用 OFFSET | 加组合索引,改键集分页 |
| 第一次执行时报错找不到游标 | sync_state表为空 | 查询 SQL 用COALESCE兜底默认时间 |
| 目标端出现源端已删除的记录 | 源端物理删除,增量感知不到 | 建立对账周期或引入 CDC |
| n8n 执行节点内存暴涨 | 一次性拉取数据量过大 | 分批拉取,循环内控制批次大小 |
| 同一条记录重复写入 | 分页重叠或网络重试 | 目标端统一使用幂等 upsert |
| 时间字段时而多 8 小时时而正常 | 时区处理不统一 | 全链路统一使用 UTC,字段名标明时区信息 |
6.2 几条必须写在工作流注释区里的规矩
工作流搭建好以后,真正维护的人可能不是你自己。所以我非常建议团队里约定几条同步工作流的硬规矩,最好直接写成注释贴在工作流里:
第一,水位线只升不降。任何情况下,都不允许直接把sync_state里的游标值改小,哪怕你怀疑漏了数据,正确做法也是先跑一次临时的补数流程,确认补完之后再把游标推进到 当前实际边界。
第二,目标端所有写入都必须是幂等的。无论增量逻辑写得再好,重复执行都应该产生和单次执行一样的结果。如果哪天下线了一个目标端表,重构时也要重新把 upsert 逻辑补上,不要偷懒改成普通的 insert。
第三,所有源端时间字段在清洗节点统一格式。这个规矩主要是为了防止后来者看图说话,不同的源端可能给datetime、timestamp、字符串类型,归一化之后,后续所有节点处理起来都一个套路,不会因为字段类型不同而分叉出多套逻辑。
第四,同步工作流一定不能只有成功路径。哪怕是一个简单流程,也要留出失败分支,把错误写进日志并发通知。没有异常感知的增量同步,生产环境里就是在裸奔。
6.3 最后再分享一个小技巧
我在项目里后期经常被问到:“为什么同步已经跑完了,数据量对比还是对不上?”排查到最后,十次里有八次是目标端的唯一键和源端不一致造成的。源端业务表的主键可能是id,但同步链路里其他模块引用的可能是business_id或者一个复合键。建议在写目标表时,不仅用主键做 upsert 冲突检测,还要把源端的原始唯一键字段原样保留下来,作为对账用的参考列。这样即使同步流程已经上线很久,你依然能随时通过字段对比定位到问题。
另外,我自己一般会给同步状态表加一个last_max_id字段,即使当前策略用时间戳,这个字段也可以顺手记录一下。将来如果你决定从时间戳切换成自增 ID 增量,或者反过来做兼容,这个字段就不需要回填历史数据,算是给未来的自己留了一条后路。
这个同步工作流从设计到现在已经稳定跑了几个月,期间经历过大促期的数据洪峰,也遇到过后端临时换表结构的突发调整。回过头来看,增量同步的核心其实不在于用了多少花哨的节点和技巧,而在于水位线管理是否严谨、目标端写入是否幂等、失败路径是否有感知。只要这三件事想透了,哪怕以后换成完全不同的数据源,也只需要改一个查询节点,整个骨架依然能继续用下去。