Flink CDC是这两年很火的数据同步工具,能实时捕获数据库变更,同步到数据湖或数仓。它基于Debezium,通过读取数据库的binlog来捕获增删改操作,延迟低,不需要轮询,对源库影响小。
我们团队最近做了一个旧系统到新系统的迁移项目。旧系统用的是单体架构,数据库是MySQL,数据量大概几亿条。新系统是微服务架构,数据库分了好几个,还需要把数据同步到数据湖做分析。我们用Flink CDC做数据实时同步,把旧库的数据实时同步到新库和数据湖,实现平滑迁移。
整个项目做了两个多月,踩了很多坑,也总结了很多经验。今天分享这次Flink CDC迁移的实战经验,包括架构设计、配置调优、遇到的问题和解决方案,以及迁移过程中的注意事项。希望能帮到正在做类似项目的朋友。
一、为什么选Flink CDC
在说实战之前,先说说为什么选Flink CDC,而不是其他工具。
我们调研了几种数据同步方案:
第一, Canal。阿里开源的MySQL binlog解析工具,很成熟,国内用得多。但Canal主要是解析binlog,数据同步需要自己写客户端,和下游系统集成要自己开发。而且Canal的HA和容错比较弱,需要自己维护。
第二,Debezium。Flink CDC底层就是用的Debezium。Debezium本身是Kafka Connect的连接器,需要配合Kafka使用,架构比较重,要维护Kafka集群。
第三,DataX/Canal+Kafka+Flink。传统的方案,用Canal读binlog写到Kafka,再用Flink消费Kafka写到下游。这个方案成熟,但组件多,运维复杂,链路长延迟高。
第四,Flink CDC。Flink CDC把binlog解析和流处理合在一起,直接读MySQL binlog,在Flink里做转换,然后写到下游。不需要Kafka做中转,架构简单,延迟低,而且能利用Flink的流处理能力做数据转换和清洗。
综合考虑,我们选了Flink CDC。主要原因是:架构简单,不需要Kafka中转;延迟低,秒级同步;能在Flink里做数据转换和清洗,灵活;Flink的Exactly-Once语义保证数据不丢不重;社区活跃,版本更新快。
我们用的是Flink 1.13加Flink CDC 2.0,这个版本支持增量快照,不需要锁表,对源库影响小。
二、整体架构设计
我们的迁移架构分三层:源库层、同步层、目标层。
源库层:旧系统的MySQL数据库,主从架构,主库写,从库读。Flink CDC连接从库读binlog,减少对主库的影响。
同步层:Flink CDC作业,部署在YARN集群上。每个作业负责同步几张表,读binlog,做数据转换(字段映射、类型转换、数据清洗),然后写到目标库。
目标层:新系统的MySQL数据库(分库分表),还有数据湖(Delta Lake on S3)。数据实时写入新库,同时双写到数据湖做分析。
迁移分两个阶段:
第一阶段是全量同步。Flink CDC先做全量快照,把历史数据全部同步到目标库。这个阶段用Flink CDC的增量快照功能,不需要锁表,边读全量边追binlog。
第二阶段是增量同步。全量同步完之后,Flink CDC继续读binlog,实时同步增量数据。这个阶段旧系统和新系统双写,数据保持一致,然后逐步把流量切到新系统,最后下线旧系统。
为了保证数据安全,我们还做了数据校验:每天定时对比源库和目标库的数据,检查有没有不一致的情况。
三、Flink CDC的配置和调优
下面说说Flink CDC的具体配置和调优,这部分是实战的重点。
1. 增量快照配置
Flink CDC 2.0支持增量快照(Incremental Snapshot),不需要锁表,这对生产环境很重要。配置很简单,设置scan.incremental.snapshot.enabled=true就行。
增量快照的原理是:先把表分成多个chunk,每个chunk读一部分数据,读的时候记录binlog位置,读完chunk后追平binlog,然后读下一个chunk。这样不需要锁表,也能保证全量数据和增量数据的一致性。
要注意的是,增量快照需要表有主键,没有主键的表不支持。我们有几张表没有主键,后来加了主键或者用唯一键代替。
2. 并行度设置
Flink CDC的源端并行度要根据表的数据量和binlog产生速度来设。全量同步阶段并行度可以设高一点,加快全量同步速度。增量同步阶段并行度不需要太高,因为binlog是单线程读的,太高也没用。
我们的经验是:全量阶段源端并行度设为4到8,根据表的数据量调。增量阶段源端并行度设为1,因为binlog只能单线程读。下游写入的并行度可以设高一点,比如8到16,提高写入速度。
3. 检查点配置
检查点(Checkpoint)是Flink容错的关键。我们设置检查点间隔为1分钟,检查点超时时间设为10分钟,最小间隔30秒。检查点存在S3上,不存本地,防止节点故障丢失。
要注意的是,Flink CDC的检查点会保存binlog位置,故障恢复后从上次的binlog位置继续读,不会丢数据。但如果检查点间隔太长,故障恢复后要追很多binlog,延迟会变大。1分钟是比较合适的间隔。
4. 背压和流量控制
全量同步阶段,源端读数据很快,下游可能跟不上,产生背压。Flink会自动处理背压,但如果背压严重,会影响checkpoint,导致checkpoint超时。
我们的做法是:全量阶段限制源端的读取速度,用scan.snapshot.fetch.size控制每次fetch的大小,设为1024。下游用批量写入,攒一批再写,提高效率。这样背压就不严重了,checkpoint也能正常完成。
5. 时区和类型转换
MySQL的datetime类型和Flink的timestamp类型有时区问题,容易差8个小时。我们在作业里统一设置时区为Asia/Shanghai,用Table API的时候做了时区转换。
还有decimal类型、json类型、text类型,都要注意转换。我们遇到过decimal精度丢失的问题,后来在Flink SQL里显式指定了精度和小数位才解决。
四、遇到的问题和解决方案
整个迁移过程中遇到了不少问题,挑几个典型的说说。
问题一:全量同步太慢
最开始全量同步一张几千万的表,跑了好几个小时还没完。排查发现是源端并行度设成了1,单线程读全量。后来开启增量快照,源端并行度设为8,全量速度提升了好几倍,一个多小时就同步完了。
还有一个原因是下游写入慢。我们最开始是单条insert,后来改成批量写入,攒500条写一次,速度提升了很多。
问题二:binlog位点丢失,作业重启后重复同步
有一次Flink作业因为集群故障重启,结果从很早的binlog位置开始读,重复同步了很多数据。排查发现是checkpoint没成功,保存的binlog位置太旧。
解决方法:调大checkpoint超时时间,确保checkpoint能成功。增加checkpoint的保留数量,保留最近的10个checkpoint。作业重启的时候指定从最新的savepoint启动,不要随便从最早的位置开始。
还有就是目标库的写入要做幂等处理,用INSERT ... ON DUPLICATE KEY UPDATE,即使重复同步也不会产生重复数据。
问题三:大事务导致binlog读取卡住
有一次旧系统做了一个大事务,更新了几百万条数据,binlog一下子产生了很多。Flink CDC读这个大事务的时候卡住了,延迟越来越大。
原因是Flink CDC处理大事务的时候,要把整个事务的事件放在内存里,等事务提交后才往下发。事务太大,内存不够,就卡住了。
解决方法:调大TaskManager的内存,给Flink CDC作业更多内存。和业务方沟通,避免一次性执行太大的事务,分批执行。如果实在有大事务,可以临时把作业暂停,等大事务执行完再恢复。
问题四:DDL变更导致作业失败
迁移过程中,旧系统的表结构变了,加了一个字段。Flink CDC作业因为Schema不匹配直接失败了。
解决方法:Flink CDC支持Schema Evolution,可以配置在遇到DDL时的行为。我们用的是兼容模式,遇到加字段的DDL能自动适应。但如果是删字段或者改类型,还是会失败。所以迁移期间要和业务方约定,尽量不做DDL变更,要变更的话提前通知,先更新Flink作业的Schema再变更。
问题五:数据不一致
迁移过程中我们做数据校验,发现有几张表的数据不一致,源库有的数据目标库没有。排查发现是这些表在全量同步期间有数据删除,Flink CDC的增量快照在某些边界情况下会漏掉删除事件。
解决方法:全量同步完成后,做一次全量数据校验,把不一致的数据补同步。增量同步阶段,每天做一次增量校验,对比最近一天的数据。发现不一致及时修复。我们写了一个数据校验工具,自动对比源库和目标库,生成差异报告。
五、迁移过程中的注意事项
最后总结一些迁移过程中的注意事项,都是血泪教训。
第一,迁移前做好数据备份。 不管同步工具多可靠,迁移前一定要备份源库数据。万一出问题,能从备份恢复。我们迁移前做了全量备份,还做了一次演练,确保备份能恢复。
第二,先在测试环境验证。 不要一上来就在生产环境跑Flink CDC同步。先在测试环境用生产数据的副本测试,验证同步的正确性、性能、容错性,没问题了再上生产。我们在测试环境跑了两周,发现并解决了很多问题。
第三,双写期间做好数据校验。 旧系统和新系统双写期间,一定要持续做数据校验,确保两边数据一致。不要等切流量了才发现数据不一致,那就晚了。我们每天做校验,持续了一个月,确认数据完全一致才切流量。
第四,灰度切流量。 不要一下子把所有流量切到新系统。先切10%,观察一段时间没问题,再切30%、50%、100%。每切一次都要监控系统状态和数据一致性,有问题及时回滚。我们灰度切流用了一周,很平稳。
第五,做好监控和告警。 Flink CDC作业的监控很重要,要监控作业状态、checkpoint状态、同步延迟、数据量。同步延迟超过阈值要告警,作业失败要告警。我们用Prometheus加Grafana做监控,关键指标设了告警,出问题能及时发现。
第六,回滚方案。 迁移一定要有回滚方案,万一新系统出问题,能快速切回旧系统。我们的回滚方案是:流量切回旧系统,Flink CDC反向同步(新库到旧库),确保数据不丢。虽然最后没用到,但准备了心里踏实。
六、写在最后
以上就是我们用Flink CDC做系统迁移的实战经验。从架构设计、配置调优,到遇到的问题和注意事项,尽量写得详细,希望能帮到大家。
总的来说,Flink CDC是一个非常优秀的数据同步工具,架构简单,延迟低,功能强大,适合做数据实时同步和系统迁移。但它也不是银弹,在大事务、DDL变更、数据一致性等方面还是有坑,需要在实践中不断摸索和优化。
2021年了,Flink CDC发展很快,2.0版本的增量快照功能解决了很多1.0的问题,社区也很活跃。如果你正在做数据同步或者系统迁移,Flink CDC是一个值得考虑的方案。
最后提醒一句,数据迁移是高风险操作,一定要谨慎。做好备份,做好测试,做好校验,做好监控,做好回滚方案,确保万无一失。数据安全永远是第一位的。
祝大家的迁移项目都能顺利完成,数据零丢失,系统平稳过渡。
评论(0)
暂无评论,快来抢沙发~
评论功能仅对会员开放,请先登录
登录