这个项目是去年我接手的一个数据同步平台,一条从订单库、CRM库、财务系统到统一汇总库的ETL链路。上线前没人觉得它复杂,就是"定时把数据搬过来",可真到做增量抽取、字段映射、幂等写入、失败重试这些环节时,我发现自己低估了它。这篇记录我尽量把决策过程、踩坑链路和能复用的经验写清楚,给准备做数据同步、报表底表或者内部系统对接的朋友留一份真实参考。
1. 这个项目为什么存在:一次数据孤岛的集中治理
1.1 三个业务系统里的同一份订单
先说背景。我们公司内部有独立的订单管理、CRM和财务核算系统,各自维护着自己的数据库。同一个客户、同一笔订单,在三个系统里用着完全不同的编号规则和字段命名。业务方要一份"全链路销售报表"时,只能把三个库的数据导出到Excel,再由报表专员手工关联清洗。这个流程每个月损耗大量人工,而且错漏率高。
我接手时,项目目标是把这些数据定期汇总到一个独立的PostgreSQL汇总库,做成统一宽表,供报表平台直接查询。听起来简单,但我给这个项目定了一个前提:尽量不改动业务系统,也不能要求业务方DBA开放过多权限。这意味着所有抽取逻辑都要放在一个独立服务里,从各个业务库读取数据,再做清洗转换。
1.2 需求评审时差点漏掉的两个边界场景
评审会上,业务方提的核心需求只有一句话:"每天能看到准确的经营数据。"这句话隐藏了很多坑,最典型的是两个边界场景。
第一个是"跨天数据漂移"。比如订单在昨天晚上创建,但支付回调今天凌晨才到达。如果我们的同步任务以"业务日期"为准,那么今天拉取的增量就漏掉了昨天那条订单的最终状态。这事如果不提前约定,等上线后业务拿着日报来对账一定会吵起来。我们最终的约定是:所有表统一以"最后更新时间"作为增量字段,而不是业务发生时间。
第二个是"已删除数据的处理"。业务系统允许作废订单,但有时直接物理删除,这会造成汇总库此前同步的记录永远留在宽表里。我们排查后确认,只有极少数表有软删除标记,其他表可能被物理删除。为了控制复杂度,第一版只同步"插入+更新",不对物理删除做额外处理,但要把这个限制明确写进需求评审结论里,否则后期一定会被当Bug追着问。
1.3 项目范围线是怎么画的
需求评审后,我把范围收敛成三块:
一是数据抽取层,负责从各业务库读取增量数据;二是清洗映射层,负责字段类型转换、单位换算、状态枚举映射;三是任务调度层,负责定时触发、依赖编排和失败告警。
我刻意没有做数据回写、实时接口和自助取数平台,因为这些需求一旦展开,项目周期会膨胀到无法交付。这种减法越早做越好,后期业务方加需求时,我可以明确说:回写功能需要单独的变更评估,而不是临时塞进同步任务里。
2. 技术选型取舍:为什么最终没有引入 Airflow 和 Canal
2.1 先说结论:团队、数据量、运维成本决定一切
技术评审时,摆在桌上的是三条路线:引入Apache Airflow做调度加ETL编排;引入Canal订阅MySQL binlog做实时增量;或者用Python自研一个轻量同步框架。
我最后选了自研,原因很现实:团队没有专门的运维岗,Airflow虽然功能强,但对基础组件依赖多,部署后还要专门维护元数据库和worker进程;而Canal的实时同步能力当时用不上,业务侧没有"秒级一致"的要求,而且binlog订阅需要业务库开启相关参数,协调成本高。数据量是另一个重要因素——单表日增量几百上千行,全量也就几十万行,这个量级用定时轮询时间戳就够了。
| 方案 | 部署成本 | 实时性 | 对业务库影响 | 团队维护难度 |
|---|---|---|---|---|
| Airflow | 高,多组件 | 好 | 需要额外账号和网络策略 | 高 |
| Canal | 高,需开启binlog | 秒级 | 需要业务库配合开启参数 | 中 |
| 自研定时任务 | 低,单服务 | 分钟级 | 只做普通查询 | 低 |
2.2 自研调度与任务框架的模块轮廓
我搭的同步服务大体分四层,第一层是调度入口,当时直接用了APScheduler的CronTrigger,一个任务配置对应一个定时表达式;第二层是抽取器,每个数据源定义一个Reader对象,负责拼接查询SQL并执行;第三层是转换器,把源字段名映射到目标字段名,做类型强制转换和默认值填充;第四层是写入器,统一走PostgreSQL的MERGE语法做幂等写入。
整个服务逻辑不复杂,但有一个关键设计:所有任务配置都写在JSON文件里,包括数据源连接、源表名、目标表名、增量字段、字段映射关系、批大小参数。新增一张表的同步任务,只需要添加一段配置,不用改代码。这让后来接手的人不需要理解Index服务内部的Python代码就能扩展同步逻辑。
2.3 我坚持的三个设计原则
第一个原则是"任务必须可重跑"。数据同步任务失败率天然高,网络抖动、锁竞争、超时都可能中断,所以每一次运行都必须是幂等的,重跑不会产生重复数据。
第二个原则是"任何状态都要可观测"。每轮任务开始和结束时都要记录start_time、end_time、源行数、目标行数、耗时,写入一张任务运行日志表。这个表后来在排障时起了关键作用,很多问题不需要去看业务表,只看日志表就能定位。
第三个原则是"查询既要快又要轻"。同步任务用的查询大部分会命中主键索引或更新时间索引,但业务系统有时会加只读从库作为数据源,避免主库压力过大。这一点要提前跟DBA确认,否则白天跑同步时,几张大表的范围查询能把主库CPU拉高。
3. 核心模块落地:从增量同步到幂等写入的每个细节
3.1 增量同步:基于时间戳的边界条件处理
时间戳增量是最常见的方案,真正的坑在边界。我的抽取SQL最早写的是WHERE update_time > last_sync_time,但很快发现两个问题。
一是"边界时间重叠"问题。如果同步任务在10:00:00.500开始抓取,最后一条数据的update_time恰好是10:00:00.600,那么下一次同步从10:00:00.500开始,这条数据就会被再抓一次。这其实不严重,因为写入端做幂等就能扛住,但会让任务运行日志里的"源行数"每次都包含一部分重复行,对排查有干扰。所以我在配置里加了一个可选参数overlap_seconds,默认向后偏移2秒,也能覆盖恰好跨越秒边界的更新。
二是"时间字段精度"问题。源库有些表update_time是datetime(0),存储精确到秒。如果业务在500毫秒内对同一行做了两次更新,第二次更新的时间戳和第一次完全相同,那么基于秒级时间戳的增量查询就会永久漏掉这条记录。这个问题上线后真实发生过,后面章节我会单独复盘。
最终我采用的增量抽取SQL类似这样:
SELECT * FROM source_table WHERE update_time > :last_sync_time AND update_time <= :current_sync_time ORDER BY id;为什么要加<= :current_sync_time上限?因为一旦任务运行时间较长,运行期间新插入的数据的更新时间会超过任务启动时间,如果不对本轮查询设上限,下一轮同步就从启动时间扫起,很容易造成大量重复或数据混乱。加上限之后,每一轮数据范围都是左开右闭区间,逻辑边界非常清晰。
3.2 字段映射层:用 JSON 配置代替硬编码清洗
清洗逻辑如果硬编码在代码里,每加一张表都要改代码发版,很快会失控。我用了一个映射配置文件来描述源字段到目标字段的关系,基本结构是这样:
{ "table_name": "order_wide", "source_table": "orders", "incremental_field": "last_modified_time", "primary_key": "order_id", "column_mapping": [ {"source": "order_no", "target": "order_code", "type": "string", "required": true}, {"source": "amount", "target": "order_amount", "type": "decimal(10,2)", "default": "0.00"}, {"source": "status", "target": "status_code", "type": "mapping", "value_map": {"1": "created", "2": "paid", "3": "closed"}} ] }其中"type"字段支持string、integer、decimal、datetime、mapping等几种。最实用的是mapping类型,用来做枚举值转换,比如业务库的订单状态是1、2、3,目标表要存created、paid、closed。这个配置化思路在后续新增表时节省了大量沟通成本。
转换逻辑的核心难点在于"脏数据不能中断整个任务"。我做的处理是:单行转换失败时记录错误详情到一张etl_error_log表,把这行跳过,而不是整体任务失败。宁可让几张表少几条数据并告警出来,也不能因为一条坏数据让所有表停摆。
3.3 写入层:唯一索引、MERGE 语句与幂等性
写入层的设计是整套系统的安全底线。目标表所有业务主键都要建唯一索引,这一点没有商量。即使源表没有主键,也要选择业务上能唯一标识行的字段组合作为唯一键。
写入时没有使用"先查再决定插入或更新"的逻辑,这样并发会有竞态问题。统一使用PostgreSQL的INSERT ... ON CONFLICT DO UPDATE,一条SQL完成更新或插入,天然保证幂等。
INSERT INTO order_wide(order_id, order_no, order_amount, status, update_time) VALUES (:order_id, :order_no, :order_amount, :status, :update_time) ON CONFLICT (order_id) DO UPDATE SET order_no = EXCLUDED.order_no, order_amount = EXCLUDED.order_amount, status = EXCLUDED.status, update_time = EXCLUDED.update_time;这里有一个细节值得单独说:要不要把更新时间戳字段一起强制覆盖。一种做法是只更新业务字段,保留目标表的首次创建时间;另一种是无条件覆盖所有字段。我选了后者,因为同步任务的目的就是让目标表尽可能等同源表,覆盖所有字段不会产生业务歧义。
3.4 任务编排:依赖关系、失败重试与断点恢复
同步任务不是互相独立的。比如订单宽表依赖订单基础表和客户映射表,两者都完成后才能触发。我第一版用APScheduler自带功能实现,后来发现任务多了后依赖关系开始混乱,就引入了一个非常轻量的方案:把任务依赖配置写在JSON里,同步服务在每轮调度时先检查前置任务的最近一次运行结果,如果前置任务失败或者未运行,则跳过本轮。
这个方案后来被证明够用,但没有解决一个重要的"断点恢复"问题。比如某个同步任务在凌晨4点运行到一半,数据库连接超时异常退出。此时已经读取了部分源数据且写入了一部分目标表,重启后若从头重跑,那数据还好;如果任务配置里设置了"只抽取最近10分钟增量",就可能漏掉连接断开前的数据。我的处理方式是:给每个任务增加一个last_processed_id参数,在每次写入一批后持久化到任务状态表,重新运行时从这个ID之后继续。这个设计是运行两个月后最实用的一个补丁。
4. 上线前后实录:四个直接把系统拖住的真实故障
4.1 主键重复:不是数据错了,是写入逻辑有竞态
运行第三天,告警群突然弹出几百条"唯一索引冲突"日志。我当时第一反应是源数据有问题,导出源表检查发现同一订单确实出现了两行,但ID相同,其他字段不完全一致。
排查发现问题出在我自己的写入逻辑上。我在一个位置为了省事,写了"先按订单号查目标表,没查到再插入",查到就更新。两个同步实例同时在运行(实际是一次手动补数任务和定时任务重叠了),它们都查出"订单不存在",然后同时执行插入,其中一个在唯一索引处碰壁。
这次故障后我把写入全部统一改为ON CONFLICT DO UPDATE,不再使用"先查再写",并且没有在应用层加分布式锁。这个问题的教训是:应用层任何形式的"检查后执行"都有竞态窗口,真正的保险永远是数据库的唯一约束和原子upsert。
4.2 DATETIME(0) 丢精度导致漏数
上线第一周后,对账时发现客户表中有一行记录很久没更新,却是最近才被业务方修改过的。我把源表和目标表的同一行逐字段对比,发现源表的确更新过,但增量查询没有抓到。
导出源表数据后定位到了原因:源表update_time是datetime(0),只精确到秒。业务在某个秒级窗口内对同一行做了两次更新,第二次更新的时间戳和第一次一样。上一轮同步已经用这个时间戳做了边界,下一轮参数是last_sync_time + overlap,因为这个时间戳没有变化,所以查询永远找不到这一行。
这个问题的解决有两个方向:一是要求业务库把时间字段改造成datetime(6)并给所有更新应用精确的当前时间;二是把增量字段换成自增主键。主键方案最稳,但源库有些表的主键不是自增数字。最终大部分表改成截取自增ID做增量标记,少数没有自增ID的表靠业务侧配合在应用层更新时间戳精度。这类源库表结构问题必须在项目开始前列清单,逐表确认增量字段精度,绝不能默认所有表都一样。
4.3 全量增量交叠期把订单表同步跑重
项目中期有一张新表需要先做全量数据初始化,再做增量同步。操作时我先启动了一个全量同步任务,又立刻启动了增量同步任务,两个任务同时写一张目标表。结果全量任务还在跑,增量任务又写入了几条之前全量已经覆盖的记录,虽然ON CONFLICT UPDATE没有导致重复主键,但那些"先全量后增量"的字段被全量任务的老快照覆盖回去了,产生了短暂的数据回退。
这类问题最有效的规避方式是给每个任务增加"运行互斥锁",同一张目标表同一时间只允许一个写入任务运行,我直接在数据库里建了一张task_lock表,任务启动时插入一行记录,结束后删除,同一张目标表第二个任务会因为拿不到锁而等待或跳过。此后这个故障再没出现过。
4.4 失败重试风暴:十分钟打满连接池
默认失败重试策略我设成了每10分钟重试一次,最多重试5次。某天业务库做维护,连接被拒绝,同步任务开始密集重试。每个任务重试时要新建数据库连接池连接,多个表一起触发后直接把源库连接数打满,导致正常业务系统也开始报连接超时。
这次事故让我把重试策略彻底改掉。首先,短时错误(连接被拒绝)走"指数退避+抖动",首次等待30秒,之后线性增长到300秒封顶;其次,不依赖应用自身的定时重试,而是让任务失败后进入待重试队列,由统一的调度器控制重试节奏;最后,同一时间只允许一定数量的任务实例重跑,而不是放所有任务一起撞上去。这里我给待重试任务维护一个信号量,默认并发上限为3,防止数据库一恢复就被同步任务再次拖垮。
5. 复盘:哪些决策省了大事,哪些地方下次绝不这么做
5.1 写进后续项目检查清单的五件事
这个项目做完后,我把几条经验沉淀成了自己后续项目的强制检查项,分享给你们参考。
第一,配置驱动的同步逻辑远优于代码硬编码。新增表的维护成本差了一个量级;第二,目标表必定要有唯一索引,写入必须用原子upsert,这条已经写进了我的代码规范;第三,增量字段必须确认精度和唯一性,不能默认时间戳一定靠谱;第四,每个任务必须有可观测的运行日志,包括任务开始时间、结束时间、抽取行数、写入行数、失败行数、耗时,这些数据既用来排查问题,也能用来做任务趋势分析;第五,所有外部依赖(业务库连接、网络、磁盘)都需要在任务异常时优雅降级,不能让任务是整个系统里最不稳定的环节。
5.2 这次没做、下次一定会做的三样东西
一是"源表结构变更的自动感知"。项目期间有一次源表加列导致映射配置报错,当时靠告警才发现。如果做一个"对比源表实际列与配置列"的启动自检,问题能在同步开始之前暴露。二是"数据质量校验规则引擎"。现在只能对账行数和主键,对金额、状态、关系完整性还没有自动校验。如果能把"订单总金额=订单明细金额之和"这类规则配置到系统里,业务信任度会提高很多。三是"按域分环境的分层部署"。测试环境直接复用生产任务配置,容易出现测试任务把生产目标表误写的问题,下一版必须拆分环境。
5.3 最后一点个人体会
做这类内部同步系统,最大的敌人往往不是技术,而是"默认一切正常"的心态。时间戳一定可靠吗?表结构一定不变吗?任务一定不重叠吗?我在这个项目里学的每一条教训,几乎都源于对某个"默认情况"的过早信任。
数据同步不是什么炫酷的架构,它更像是修管道——大部分时间不被人注意,一旦漏了水就是事故。把边界想清楚,把日志做扎实,把重跑做成常态,这套系统就能安安稳稳地运行下去。如果这个记录对你的同步项目有一点点参考价值,那这篇文章就没白写。