很多企业建设实时数仓、经营驾驶舱或风险预警系统时,第一反应是把数据同步频率从“每天一次”调成“每分钟一次”。
但真正的实时同步,并不是让任务跑得更勤。
源数据库持续发生新增、修改和删除,全量迁移期间仍然有新数据写入,网络可能中断,目标库可能限流,表结构也可能随时调整。任何一个环节没有处理好,都可能造成数据丢失、重复、错序或长期不一致。
为方便大家系统梳理数据采集、数据集成、数据治理和分析应用的完整建设路径,我整理了一份数字化全流程资料包,需要自取:https://s.fanruan.com/tyac0(复制到浏览器)
一套可靠的实时同步方案,至少要回答六个问题:
变化如何捕获、历史数据如何迁移、全量与增量如何衔接、目标库如何写入、异常如何恢复、结果如何校验。
一、先定义实时目标,不要一上来就选工具
“实时”没有统一标准。
交易反欺诈可能要求秒级甚至更低延迟;设备告警可能要求十秒级;经营看板允许几分钟延迟;财务分析则可能只需要每小时更新。
因此,在设计同步方案前,首先要确定四类指标:
正常情况下允许多少延迟;
业务高峰期最多允许积压多久;
任务中断后需要多长时间恢复;
是否允许短时间内出现重复数据。
同时,还要明确同步范围。
有些表只需要同步新增数据,例如日志、流水和设备采集记录;有些表则必须完整捕获新增、修改和删除,例如订单、客户、合同和库存数据。
延迟越低,对源库、网络、目标端和运维能力的要求就越高。
所以,实时同步并不是越快越好,而是要让同步时效与业务价值相匹配。
在方案规划阶段,我通常会先把数据源、目标端、同步表、业务主键和更新方式放到FineDataLink中统一梳理,再按照业务需求划分实时链路、分钟级链路和定时批量链路。这样做的重点不是“少写几段代码”,而是先建立一张完整的数据流转地图,避免同一张表被重复抽取、不同系统使用不同更新规则。
二、源端采集:关键不是读到数据,而是识别变化
实时同步的第一步,是判断源数据库中发生了什么变化。
1、基于数据库日志的CDC
CDC,即变化数据捕获。它通常通过读取数据库事务日志,识别INSERT、UPDATE和DELETE事件。
相比反复扫描业务表,CDC具有几个明显优势:
对源数据库影响相对较小;
能够识别新增、更新和物理删除;
可以保留数据变化的大致顺序;
适合数据量大、时效要求高的场景。
但CDC并不是配置完成就可以长期稳定运行。
首先要确认数据库是否开启完整日志,日志保留周期是否足够,同步账号是否具备读取权限。
假设同步任务中断两天,而数据库日志只保留一天,那么即使系统支持断点续传,也无法从原来的位置恢复,只能重新进行全量同步。
因此,日志保留周期必须覆盖系统可能出现的最大故障恢复窗口。
其次,还要关注大事务问题。
如果业务系统一次性更新几十万条数据,数据库可能在事务提交后集中释放大量变化记录。同步链路会在短时间内收到大量事件,造成消息积压和目标端写入压力。
所以,CDC只解决了“变化如何发现”,并不代表这些变化一定能够及时写入目标库。
2、基于更新时间戳抽取
时间戳增量通常读取“更新时间大于上次同步时间”的数据。
这种方式实现简单,但边界问题非常明显。
例如,上次同步到10点整,下次查询条件设置为“更新时间大于10点”,那么恰好在10点整更新的数据就可能被漏掉。
更稳妥的方式有两种:
一种是使用“更新时间+主键”作为复合游标;
另一种是设置回溯窗口,每次多读取前几分钟的数据,再通过主键和版本号去重。
增量同步宁可重复读取少量数据,也不能为了避免重复而留下漏数风险。
时间戳方案还有一个天然缺陷:无法识别物理删除。对于存在删除操作的业务表,需要改成逻辑删除、增加删除日志表,或者直接采用CDC方案。
3、基于递增主键抽取
如果数据只新增、不修改,可以记录上次同步的最大ID,下次只读取更大的数据。
这种方案适合日志、流水等严格追加型数据,却不适合订单、库存和客户信息。
因为最大ID只能说明新增到了哪里,不能判断历史记录是否被修改,也无法发现删除操作。
如果业务存在历史补录、主键跳号、分库合并,或者ID生成顺序与事务提交顺序不一致,最大ID甚至不能代表之前的数据已经全部完成同步。
三、全量初始化:难点在存量与增量的交界
目标库第一次建设时,通常需要先同步历史数据,再进入实时增量阶段。
真正容易出问题的,不是全量怎么复制,而是全量执行期间,源数据库仍然在持续变化。
假设上午10点开始复制订单表,中午12点完成。
如果等到12点以后再启动增量任务,那么10点到12点之间发生的新增、修改和删除就可能丢失。
一套更稳妥的流程通常包括六步:
在全量开始前记录数据库日志位点;
启动增量变化捕获,并将变化暂时保存;
按照主键范围或时间分区读取历史数据;
全量完成后,从已记录位点回放增量;
根据主键和版本号处理重复与覆盖;
完成数据校验后,切换为持续增量同步。
这里必须分清两个概念:
全量数据代表读取时看到的数据状态,增量日志代表之后发生的数据变化过程。
例如,一条订单在全量读取时是“待支付”,随后增量日志中出现“已支付”和“已发货”。即使全量批次晚于增量事件到达目标端,最终结果也必须保持为“已发货”,不能被旧状态重新覆盖。
进入这一阶段后,FineDataLink的实时管道可以把历史存量同步和后续增量捕获放在同一条任务链路中管理。存量负责建立目标库基线,增量负责持续追踪变化,任务切换不再完全依赖人工记录日志位置和临时修改脚本,能够降低全量与增量交界处出现数据空档的概率。
全量同步还不能简单地一次性读取整张大表。
对于上亿行数据,应按照主键区间、时间分区或者业务区域分批抽取,并限制并发数。否则,全量任务可能长时间占用数据库连接、磁盘I/O和网络带宽,影响正常业务交易。
四、传输层设计:既要承接高峰,也要控制顺序
数据变化量较小时,可以从源端直接写入目标库。
但在高并发、高波动场景下,通常需要在中间增加消息队列或其他缓冲机制。
假设业务高峰期每秒产生2万条变化,而目标库每秒只能稳定写入8000条。如果没有缓冲,目标库会持续承受冲击,最终导致锁等待、连接耗尽甚至任务失败。
缓冲层主要解决三个问题。
1、削峰填谷
业务高峰期先把数据暂存在队列中,目标端按照自身能力持续消费。
2、解耦源端与目标端
目标库短暂维护或性能下降时,不需要阻塞源系统正常交易。
3、支持多个下游消费
同一份订单变化,可以同时提供给实时数仓、风控系统、搜索服务和消息通知系统。
但引入缓冲层后,也会增加新的设计问题。
同步任务不能只监控“成功”或“失败”,还要持续观察:
当前消息积压量;
最老一条消息等待了多久;
数据生产速度是否大于消费速度;
按照当前速度多久才能消化积压。
任务没有报错,不代表链路仍然实时。
如果积压量持续增加,任务虽然显示正在运行,但目标端的数据可能已经落后数小时。
此外,还要处理事件顺序。
同一订单可能依次经历“待支付—已支付—已发货”。如果不同状态被并行处理,后发生的状态可能先到,旧状态反而后写入,最终导致订单状态回退。
常见做法是按照订单ID、客户ID或设备ID进行分区,让同一个业务对象的变化进入同一处理队列。
但分区键过于集中也可能形成热点。因此,分区设计需要在顺序性和并发度之间取得平衡。
五、目标库写入:核心不是INSERT,而是幂等
很多同步任务测试时没有问题,上线后却不断出现重复数据,根本原因通常是目标端写入不具备幂等性。
所谓幂等,是指同一条数据无论被处理一次还是多次,最终结果都保持一致。
网络超时、任务重启和失败重试,都可能导致同一条变化被重复投递。
如果目标端每次都执行INSERT,就会产生重复记录;如果按照业务唯一键执行UPSERT或MERGE,重复执行也不会改变最终结果。
因此,目标表至少需要明确以下规则:
哪个字段是业务唯一键;
新增、修改和删除如何映射;
重复数据如何识别;
新旧版本如何判断;
写入失败后如何重试;
逻辑删除和物理删除如何处理。
仅有主键还不够。
假设订单版本3先到,版本2后到,如果目标端只按照主键覆盖,最终状态就会被旧版本重新写回。
更可靠的方式是在目标端保存日志序号、更新时间或业务版本号,只接受比当前版本更新的数据。
幂等解决重复,版本控制解决乱序,两者缺一不可。
到了目标端,FineDataLink不只是负责把数据写进目标表,更重要的是把业务主键、更新方式、删除策略和字段映射集中配置。相比把追加、覆盖、更新和删除规则散落在不同脚本里,统一管理更容易确认一条数据最终以什么方式落库,也能减少多人开发造成的规则不一致。
目标端还需要合理设置批量提交大小。
逐条写入延迟较低,但连接和事务开销较大;批次过大虽然吞吐量高,却容易增加锁等待,一旦失败还需要整批重试。
因此,批量大小必须根据单条数据体积、目标库性能、并发任务数和时效要求测试确定,不能直接套用固定参数。
六、如何保证数据不丢、不重、不错
在分布式同步链路中,想要实现绝对只传输一次,成本通常非常高。
更常见、也更容易落地的方案是:
至少一次投递,加目标端幂等写入。
1、不丢:目标成功后再推进位点
同步位点代表任务已经处理到哪个日志位置、时间点或消息偏移量。
如果数据还没有成功写入目标库,就提前提交位点,一旦任务中断,这部分数据就无法重新获取。
正确顺序应该是:
读取数据—写入目标—确认成功—更新位点。
如果目标端写入失败,位点不能向前推进,任务恢复后应从原位置重新读取。
真正进入生产运行后,可以把FineDataLink作为同步位点、任务状态和异常恢复的统一管理入口。链路中断后按照最近成功位置继续执行,而不是从头重跑整张大表,这样既能缩短恢复时间,也能降低大规模补数对源库和目标库造成的二次冲击。
2、不重:允许重复传输,不允许重复结果
为了避免数据丢失,系统通常需要在失败后重试。
只要允许重试,就可能出现重复投递。
因此,不能把“不重”理解成数据在链路中绝对只出现一次,而应该保证重复数据到达目标端后,不会形成重复结果。
常见措施包括:
使用业务唯一键;
采用UPSERT或MERGE;
保存事件ID;
对重复批次进行识别;
更新前比较版本号。
3、不错:控制版本和业务状态
数据同步正确,不只是字段值相同,还要符合业务逻辑。
例如:
订单不能从已完成回退到处理中;
库存不能被旧批次覆盖;
删除事件不能晚于重新新增事件并误删新记录;
同一笔支付不能被重复累计;
主表和明细表不能长期处于不一致状态。
这些问题不能只依靠数据库主键解决,还需要把业务版本、状态转换和时间顺序纳入同步规则。
七、数据校验:任务成功不等于数据可信
同步平台显示“任务成功”,只能说明程序没有报错,不能证明源端和目标端完全一致。
完整的数据校验至少应包含四层。
第一层:数量校验
比较相同时间窗口内源端和目标端的新增数、更新数、删除数和最终写入数。
但数量相同不代表数据一定正确。
源端少一条、目标端多一条,总数仍然可能相等。
第二层:主键校验
检查源端存在、目标端缺失的主键,以及目标端重复的业务主键。
对于大表,可以按照日期分区、主键范围或者哈希桶分批比较,避免全表校验占用过多资源。
第三层:字段校验
对关键字段计算汇总值、哈希值或分布情况,例如:
订单金额合计;
库存数量合计;
不同状态的数据量;
空值比例;
最大值和最小值;
关键字段哈希结果。
第四层:业务校验
即使技术字段完全一致,也不代表数据能够直接使用。
还要检查:
支付金额是否超过订单金额;
已发货订单是否缺少客户信息;
库存是否出现异常负数;
已删除合同是否仍产生后续单据;
某个时间段的数据量是否突然归零;
某类业务状态是否明显偏离历史分布。
技术校验回答“数据有没有同步一致”,业务校验回答“同步后的数据是否可信”。
校验发现差异后,还要形成处理闭环:由谁确认、是否需要补数、补数范围是什么、补数后如何复核,以及同类问题是否需要修改同步规则。
八、异常恢复:先判断异常类型,再决定是否重试
实时同步异常不能统一设置为无限重试。
不同异常应该采用不同处理方式。
1、临时异常
例如网络抖动、连接超时、数据库短暂重启。
这类异常可以自动重试,并采用逐步延长等待时间的方式,避免短时间内不断冲击系统。
2、数据异常
例如字段长度超限、非法字符、类型转换失败。
这类问题不应阻塞整条同步链路。可以将异常记录写入隔离表或异常队列,主任务继续执行,修复后再进行补写。
3、结构异常
例如源表新增字段、删除字段或修改字段类型。
新增可选字段可能可以自动兼容,但字段缩短、类型冲突、主键变化等高风险操作,应该暂停任务并及时告警。
异常恢复还要避免一个常见错误:
任务失败后直接执行全量重跑。
一张上亿行的大表如果因为最后一个批次失败就从头开始,不仅浪费资源,还可能再次冲击源库和目标库。
更合理的方式是保存批次进度和同步位点,从最近成功的位置继续执行。
结语
实时数据同步真正考验的,不是能不能把数据从源库搬到目标库,而是当数据持续变化、网络出现波动、目标端偶尔失败时,系统仍然能够恢复、校验和追踪。
一套可靠的方案应形成完整闭环:
源端准确捕获变化,
全量与增量无缝衔接,
传输层承接业务波动,
目标端实现幂等写入,
位点机制支持任务恢复,
校验机制确认最终结果。
只有把这些环节同时设计清楚,实时同步才能从一个“正在运行的任务”,真正变成企业可以长期依赖的数据基础设施。