
这个项目是去年我接手的一个数据同步平台一条从订单库、CRM库、财务系统到统一汇总库的ETL链路。上线前没人觉得它复杂就是定时把数据搬过来可真到做增量抽取、字段映射、幂等写入、失败重试这些环节时我发现自己低估了它。这篇记录我尽量把决策过程、踩坑链路和能复用的经验写清楚给准备做数据同步、报表底表或者内部系统对接的朋友留一份真实参考。1. 这个项目为什么存在一次数据孤岛的集中治理1.1 三个业务系统里的同一份订单先说背景。我们公司内部有独立的订单管理、CRM和财务核算系统各自维护着自己的数据库。同一个客户、同一笔订单在三个系统里用着完全不同的编号规则和字段命名。业务方要一份全链路销售报表时只能把三个库的数据导出到Excel再由报表专员手工关联清洗。这个流程每个月损耗大量人工而且错漏率高。我接手时项目目标是把这些数据定期汇总到一个独立的PostgreSQL汇总库做成统一宽表供报表平台直接查询。听起来简单但我给这个项目定了一个前提尽量不改动业务系统也不能要求业务方DBA开放过多权限。这意味着所有抽取逻辑都要放在一个独立服务里从各个业务库读取数据再做清洗转换。1.2 需求评审时差点漏掉的两个边界场景评审会上业务方提的核心需求只有一句话每天能看到准确的经营数据。这句话隐藏了很多坑最典型的是两个边界场景。第一个是跨天数据漂移。比如订单在昨天晚上创建但支付回调今天凌晨才到达。如果我们的同步任务以业务日期为准那么今天拉取的增量就漏掉了昨天那条订单的最终状态。这事如果不提前约定等上线后业务拿着日报来对账一定会吵起来。我们最终的约定是所有表统一以最后更新时间作为增量字段而不是业务发生时间。第二个是已删除数据的处理。业务系统允许作废订单但有时直接物理删除这会造成汇总库此前同步的记录永远留在宽表里。我们排查后确认只有极少数表有软删除标记其他表可能被物理删除。为了控制复杂度第一版只同步插入更新不对物理删除做额外处理但要把这个限制明确写进需求评审结论里否则后期一定会被当Bug追着问。1.3 项目范围线是怎么画的需求评审后我把范围收敛成三块一是数据抽取层负责从各业务库读取增量数据二是清洗映射层负责字段类型转换、单位换算、状态枚举映射三是任务调度层负责定时触发、依赖编排和失败告警。我刻意没有做数据回写、实时接口和自助取数平台因为这些需求一旦展开项目周期会膨胀到无法交付。这种减法越早做越好后期业务方加需求时我可以明确说回写功能需要单独的变更评估而不是临时塞进同步任务里。2. 技术选型取舍为什么最终没有引入 Airflow 和 Canal2.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 最后一点个人体会做这类内部同步系统最大的敌人往往不是技术而是默认一切正常的心态。时间戳一定可靠吗表结构一定不变吗任务一定不重叠吗我在这个项目里学的每一条教训几乎都源于对某个默认情况的过早信任。数据同步不是什么炫酷的架构它更像是修管道——大部分时间不被人注意一旦漏了水就是事故。把边界想清楚把日志做扎实把重跑做成常态这套系统就能安安稳稳地运行下去。如果这个记录对你的同步项目有一点点参考价值那这篇文章就没白写。