最近负责了公司实时计算平台的迁移项目,把原来基于Storm的旧实时计算系统迁移到Apache Flink上,整个过程踩了很多坑,也积累了不少经验。今天来分享一下这次迁移的实战过程,包括背景、方案、步骤、踩过的坑以及最终的效果,希望能给做类似迁移的朋友一些参考。
一、项目背景
先聊聊为什么要做这次迁移。
我们公司的实时计算系统是几年前搭建的,用的是Apache Storm,当时Storm还是很流行的流处理框架。但是随着业务的发展,数据量越来越大,实时计算的需求越来越复杂,Storm的问题也越来越明显。
Storm的主要问题包括:吞吐量不够,随着数据量增长经常出现数据积压处理不过来的情况;延迟不稳定,数据量大的时候延迟会明显升高;状态管理弱,本身没有很好的状态管理机制,需要自己实现状态的存储和管理,很麻烦也容易出问题;容错机制简单,出现故障的时候可能会丢数据或者重复处理数据,一致性保证不够;SQL支持差,业务人员不能直接写SQL做实时分析,需要开发人员写代码,效率低;社区不活跃,更新慢新特性少,而Flink的社区非常活跃发展很快。
Apache Flink作为新一代的流处理框架,在这些方面都比Storm好很多,吞吐量高、延迟低且稳定,有完善的状态管理和容错机制,支持Exactly-Once语义,支持SQL,社区活跃生态完善。所以我们决定把实时计算系统从Storm迁移到Flink。
二、迁移前的准备
在正式迁移之前,我们做了很多准备工作。
首先梳理了旧Storm系统上所有的实时计算任务,一共有30多个任务,包括实时数据清洗、实时指标统计、实时告警、实时推荐等,每个任务的逻辑、数据量、延迟要求都不一样。我们把这些任务分了类,按复杂度和重要性排序,先迁移简单的不重要的任务,积累经验,再迁移复杂的重要的任务,降低风险。
然后搭建了新的Flink集群,用的是Flink 1.11版本,集群模式用的是YARN模式,因为我们已经有Hadoop YARN集群了,复用现有资源比较方便。集群规模一开始是10个TaskManager,每个4核8G内存,后来根据任务的需求做了调整。同时也搭建了相关的配套组件,比如Kafka作为消息队列,ZooKeeper作为协调服务,HDFS作为状态后端存储,Prometheus和Grafana作为监控系统等。
我们也制定了一些技术规范,比如数据格式统一用JSON或者Avro,状态后端用RocksDB支持大状态和增量检查点,检查点间隔设置为1分钟保证数据一致性,时间语义用事件时间更准确,窗口用滚动窗口或者滑动窗口根据业务需求,所有任务都要有监控和告警。
因为团队之前主要用Storm,对Flink不熟悉,所以我们做了几次内部培训,讲了Flink的基本概念、编程模型、状态管理、容错机制、SQL等,让团队成员都对Flink有基本的了解,能上手开发。
三、迁移的整体方案
我们的迁移方案是"双跑+灰度+切流"的方式,保证迁移过程中业务不中断、数据不丢失。
具体步骤是:第一步开发Flink版本的任务,把Storm任务用Flink重新实现一遍,逻辑和原来一致;第二步双跑验证,让Storm任务和Flink任务同时运行,消费同样的数据,对比两个任务的输出结果,确保Flink任务的逻辑正确结果一致;第三步灰度切流,验证通过后先把一小部分流量切到Flink任务上,观察一段时间没问题再逐步增加流量,直到全部切到Flink;第四步下线旧任务,全部切到Flink后观察一段时间没问题就下线Storm任务,完成迁移。
这个方案的好处是风险低,业务不中断,有问题可以随时切回旧系统,保证业务稳定。
四、迁移的具体过程
接下来聊聊迁移的具体过程和一些关键步骤。
第一步是简单任务迁移。我们先选了几个简单的任务做试点,比如实时数据清洗任务,逻辑比较简单,就是从Kafka读数据,做一些清洗和转换,再写到Kafka或者数据库。用Flink实现这些任务比较简单,用DataStream API或者Table API/SQL都能实现,我们主要用DataStream API,因为更灵活。
开发完之后双跑验证,对比结果,发现Flink任务的结果和Storm任务的结果基本一致,只有少量差异,主要是时间语义的问题。Storm用的是处理时间,Flink我们用的是事件时间,所以窗口的结果有一些差异,调整了时间语义之后就一致了。
然后灰度切流,先切10%的流量观察了一天没问题,再切50%观察一天没问题,最后全部切过去,观察了一周没问题就下线了对应的Storm任务。这几个简单任务的迁移很顺利,也让团队积累了Flink开发和运维的经验,为后面复杂任务的迁移打下了基础。
第二步是中等复杂度任务迁移。然后我们迁移了一些中等复杂度的任务,比如实时指标统计任务,需要做窗口聚合、状态管理等,比简单任务复杂一些。这些任务用Flink实现也不太难,主要用到了Flink的窗口函数、状态后端等功能,Flink的这些功能都很完善,比Storm自己实现方便很多。
在迁移这些任务的过程中,我们遇到了一些问题,比如大状态的性能问题、检查点超时问题等,后面会详细说这些坑和解决方法。
第三步是复杂任务迁移。最后我们迁移了几个最复杂的任务,比如实时推荐任务、实时风控任务,这些任务逻辑复杂、状态大,对延迟和一致性要求高,迁移难度也最大。这些任务用Flink实现花了比较多的时间,也踩了很多坑,但是最终都成功迁移了,而且性能和稳定性都比原来的Storm任务好很多。
五、踩过的坑和解决方法
在迁移的过程中我们踩了很多坑,这里分享几个比较典型的坑和解决方法。
坑1:状态太大检查点超时
我们有一个实时统计任务,需要保存用户的历史行为数据,状态很大有几十G,用RocksDB状态后端,但是检查点经常超时,导致任务失败重启。
原因是检查点超时主要因为状态太大,全量检查点需要把所有状态都写到HDFS,时间很长超过了检查点的超时时间。
解决方法包括:开启增量检查点,RocksDB支持增量检查点,只把变化的部分写到HDFS,大大减少了检查点的数据量和时间;调整检查点间隔和超时时间,把检查点间隔调大到2分钟,超时时间调大到10分钟;优化状态结构,减少不必要的状态数据,比如设置状态的TTL自动清理过期的状态数据;增加TaskManager的内存和磁盘资源,让RocksDB有足够的资源。优化之后检查点就不再超时了,任务运行稳定。
坑2:数据倾斜导致个别Task处理不过来
我们有一个任务按用户ID做keyBy然后做窗口聚合,但是有一些热门用户数据量特别大,导致这些key所在的Task数据积压处理不过来,延迟很高,而其他Task很空闲。这是典型的数据倾斜问题,少数key的数据量特别大导致负载不均。
解决方法包括:对key做加盐打散,把一个大key拆成多个子key分散到不同的Task上处理,然后再做二次聚合把结果合并;开启Flink的推测执行或者本地聚合,减少数据倾斜的影响;对热门key做特殊处理,单独拉出来处理不和其他key混在一起。优化之后数据倾斜的问题得到了缓解,任务的延迟降了下来。
坑3:Kafka消费位移提交问题
迁移过程中我们遇到了Kafka消费位移提交的问题,有时候任务重启后会重复消费数据或者丢数据。主要是因为Flink的Kafka消费者位移提交和检查点的关系没搞清楚,Flink的Kafka消费者位移是在检查点完成后才提交的,这样才能保证Exactly-Once语义,如果自己手动提交位移就会和检查点冲突,导致数据重复或者丢失。
解决方法是不要手动提交Kafka位移,让Flink自动管理位移,在检查点完成后自动提交;开启检查点,设置合理的检查点间隔保证数据一致性;设置正确的消费启动模式,比如从提交的位移开始消费或者从最新开始消费,根据业务需求。优化之后位移提交就正常了,不再出现数据重复或者丢失的问题。
坑4:窗口计算结果和旧系统不一致
双跑验证的时候发现Flink任务的窗口计算结果和Storm任务的结果有差异不一致。主要是时间语义的差异,Storm任务用的是处理时间窗口按处理数据的时间来计算,而Flink任务我们用的是事件时间窗口按数据本身的时间来计算,所以结果有差异。另外还有迟到数据的处理差异,Storm任务对迟到数据没有特殊处理,而Flink的事件时间窗口有watermark和allowedLateness机制,对迟到数据的处理不同。
解决方法是统一时间语义,如果要和旧系统结果一致就也用处理时间,如果要更准确就用事件时间,但是要和业务方沟通说明差异的原因;调整watermark和allowedLateness的设置,合理处理迟到数据;双跑验证的时候考虑时间语义的差异,不要简单地对比结果,要理解差异的原因。我们最终选择了事件时间,因为更准确更合理,和业务方沟通后他们也认可事件时间的结果,虽然和旧系统有一些差异,但是更准确。
坑5:任务反压问题
有一些任务运行一段时间后出现反压问题,上游数据积压处理不过来,延迟升高。反压的原因有很多,可能是下游算子处理慢,可能是状态太大,可能是数据倾斜,可能是资源不够等。
解决方法包括:用Flink的Web UI查看反压的位置,定位是哪个算子出现反压;分析反压的原因针对性地优化,比如下游处理慢就优化代码或者增加并行度,状态太大就优化状态结构,数据倾斜就做打散等;增加TaskManager的资源或者调整并行度,让任务有足够的资源处理数据。通过这些方法我们解决了反压的问题,任务运行稳定。
六、迁移的效果
经过几个月的努力,我们把所有30多个实时计算任务都从Storm迁移到了Flink,迁移完成后效果很明显。
吞吐量提升,Flink的吞吐量比Storm高很多,同样的资源下吞吐量提升了3-5倍,原来经常数据积压的任务现在都能及时处理了,不再积压。
延迟降低且稳定,Flink的延迟比Storm低而且更稳定,原来数据量大的时候延迟会明显升高,现在延迟基本稳定在秒级甚至亚秒级,业务体验好了很多。
状态管理更完善,Flink的状态管理比Storm完善很多,不用自己实现状态的存储和管理,开发效率高了很多,而且状态的一致性和可靠性也更好。
数据一致性更好,Flink支持Exactly-Once语义,通过检查点和两阶段提交保证数据不丢不重复,一致性比Storm好很多,业务数据更准确。
SQL支持好,Flink的SQL支持很好,业务人员可以直接写SQL做实时分析,不用依赖开发人员,效率提升了很多。
运维更简单,Flink的运维比Storm简单,有完善的Web UI、监控告警工具,问题定位更方便,而且社区活跃遇到问题容易找到解决方案。
资源成本降低,因为Flink的吞吐量更高、资源利用率更高,我们减少了集群的规模,从原来的20台Storm服务器减少到12台Flink服务器,资源成本降低了约40%。
七、迁移的经验总结
最后总结一下这次迁移的一些经验。
第一,充分准备很重要。迁移之前一定要做好充分的准备,梳理旧系统的任务,搭建新集群,制定规范,培训人员等,准备充分了迁移过程才会顺利。
第二,先简单后复杂逐步迁移。不要一上来就迁移最复杂的任务,先迁移简单的任务积累经验,再迁移复杂的任务,逐步推进降低风险。
第三,双跑验证不可少。迁移过程中一定要双跑验证,让旧系统和新系统同时运行对比结果,确保新系统的逻辑正确结果一致,不要直接切流,否则出了问题很难发现。
第四,灰度切流降低风险。切流的时候要灰度逐步切,先切一小部分流量观察没问题再增加,直到全部切过去,这样出了问题影响范围小,也容易回滚。
第五,监控告警要完善。迁移过程中一定要有完善的监控和告警,监控任务的状态、延迟、吞吐量、错误等,有问题及时发现及时处理,不要等业务方反馈才知道出了问题。
第六,遇到问题不要慌,定位原因再解决。迁移过程中肯定会遇到各种问题,不要慌,先定位问题的原因再针对性地解决,Flink的社区很活跃文档也很全,大部分问题都能找到解决方案。
第七,文档和总结很重要。迁移过程中要做好文档,记录每个任务的迁移过程、遇到的问题、解决方法等,迁移完成后做总结沉淀经验,方便以后参考,也方便团队其他成员学习。
八、写在最后
好了,关于Flink实时计算迁移的实战分享就到这里。
总结一下,我们把基于Storm的旧实时计算系统迁移到了Apache Flink,采用双跑+灰度+切流的方案逐步迁移了30多个任务,过程中踩了很多坑也积累了不少经验,迁移完成后吞吐量提升了3-5倍,延迟降低且稳定,资源成本降低了约40%,效果很明显。
Flink作为新一代的流处理框架确实比Storm优秀很多,在吞吐量、延迟、状态管理、容错、SQL支持等方面都有明显优势,是实时计算的好选择。
如果你也在考虑把实时计算系统迁移到Flink,希望这篇文章能给你一些参考和启发。如果有什么问题或者不同的看法,欢迎在评论区交流讨论。
评论(0)
暂无评论,快来抢沙发~
评论功能仅对会员开放,请先登录
登录