最近负责了一个大数据迁移项目,把旧的Hadoop MapReduce系统,迁移到新的Spark系统。整个过程,踩了很多坑,也积累了很多经验。今天,分享一下这次大数据迁移的实战经验,包括迁移前的准备、迁移的步骤、数据校验、性能调优、遇到的问题和解决方案,希望对正在做大数据迁移的朋友有所帮助。

先说说项目背景。我们公司有一个大数据系统,已经运行了五六年了,用的是Hadoop MapReduce,每天处理几亿条数据,生成各种报表和指标。但随着数据量越来越大,MapReduce的性能问题越来越明显,一个任务要跑好几个小时,而且开发效率很低,写一个MapReduce任务,要写很多样板代码,维护起来很困难。

于是,公司决定,把旧的MapReduce系统,迁移到新的Spark系统。Spark的性能比MapReduce好很多,而且开发效率高,API友好,支持SQL、流处理、机器学习等,能满足我们的各种需求。

我负责这个迁移项目的技术方案和实施。整个项目,花了三个月的时间,把旧系统的几十个MapReduce任务,全部迁移到了Spark,性能提升了3-5倍,开发效率也大大提升。

今天,就把这次迁移的实战经验,分享出来。

一、迁移前的准备:不打无准备之仗

大数据迁移,是一个复杂的工程,不能上来就写代码,迁移前的准备工作非常重要。准备工作做得好,迁移过程就会顺利很多;准备工作做得不好,迁移过程中会遇到各种问题,甚至可能失败。

1. 梳理旧系统的全貌

迁移之前,首先要梳理旧系统的全貌,搞清楚旧系统有哪些任务,每个任务是干什么的,输入是什么,输出是什么,依赖关系是什么,数据量有多大,运行时间有多长。

我们花了一周的时间,把旧系统的几十个MapReduce任务,全部梳理了一遍,整理成了一个清单,包括:

  • 任务名称和功能说明
  • 输入数据源(哪些表、哪些文件)
  • 输出结果(写到哪些表、哪些文件)
  • 任务之间的依赖关系
  • 每天的数据量
  • 平均运行时间
  • 负责人和业务方

梳理完之后,我们对旧系统的全貌,就有了一个清晰的认识,知道哪些任务是核心的,哪些是次要的,哪些任务之间有依赖,哪些任务可以并行迁移。

2. 评估迁移的工作量和优先级

梳理完旧系统之后,就要评估每个任务的迁移工作量,以及迁移的优先级。不是所有任务都要同时迁移,要有计划、分批次地迁移。

我们评估每个任务的迁移工作量,主要考虑以下几个因素:

  • 任务的复杂度:逻辑简单的任务,迁移工作量小;逻辑复杂的任务,迁移工作量大。
  • 代码的质量:旧代码写得好、有注释的,迁移起来容易;旧代码写得乱、没有注释的,迁移起来困难。
  • 依赖的复杂度:依赖少的任务,迁移起来容易;依赖多的任务,迁移起来困难。
  • 业务的重要性:核心业务的任务,要优先迁移;次要业务的任务,可以后面迁移。

根据这些因素,我们把任务分成了三批:

  • 第一批:简单的、独立的、非核心的任务,用来练手,积累经验。
  • 第二批:中等复杂度的、有一定依赖的任务,占大部分。
  • 第三批:复杂的、核心的、依赖多的任务,最后迁移,确保万无一失。

这样分批次迁移,风险可控,而且可以在迁移过程中,积累经验,优化流程,后面的迁移会越来越顺利。

3. 搭建新系统的环境

迁移之前,还要搭建好新系统的环境,包括Spark集群、Hadoop集群、数据仓库、调度系统、监控系统等。

我们搭建了一套独立的Spark集群,和旧的MapReduce集群分开,互不影响。这样,迁移过程中,旧系统正常运行,新系统也在运行,两边并行,不会影响业务。

同时,我们还搭建了数据校验工具、性能监控工具、日志分析工具等,方便迁移过程中的数据校验和问题排查。

4. 制定迁移规范和标准

迁移之前,还要制定迁移规范和标准,确保所有迁移的任务,都按照统一的规范来写,避免每个人写的风格不一样,后期维护困难。

我们制定的规范包括:

  • 代码规范:命名规范、注释规范、目录结构等。
  • 开发规范:用Spark SQL还是DataFrame API,怎么处理数据,怎么处理异常等。
  • 配置规范:参数配置、资源配置、序列化配置等。
  • 测试规范:单元测试、集成测试、数据校验等。
  • 上线规范:上线流程、回滚方案、监控告警等。

制定了规范之后,所有开发人员都按照规范来写,迁移出来的代码,风格统一,质量有保障,后期维护也容易。

二、迁移的实施:分批次,稳扎稳打

准备工作做好之后,就开始正式迁移了。我们按照之前制定的计划,分三批次迁移,稳扎稳打,确保每个任务迁移成功。

1. 第一批:简单任务练手

第一批,我们选了5个简单的、独立的、非核心的任务,用来练手,积累经验。

这些任务,逻辑都比较简单,就是读取数据,做一些简单的清洗和转换,然后写入结果。我们用Spark SQL重写了这些任务,因为这些任务的逻辑,用SQL就能很方便地表达,不需要写复杂的代码。

迁移的过程中,我们遇到了一些小问题,比如:

  • 数据类型不一致:旧系统是MapReduce,自己处理数据类型,新系统用Spark SQL,对数据类型要求严格,有些字段类型不匹配,需要转换。
  • 空值处理不一致:旧系统对空值的处理,和Spark SQL不一样,需要调整。
  • 日期格式不一致:旧系统的日期格式,和Spark SQL的默认格式不一样,需要指定格式。

这些问题,都比较小,很快就解决了。第一批任务迁移完成之后,我们做了数据校验,确认新系统的结果和旧系统完全一致,性能也提升了不少。

第一批任务的成功,给了我们很大的信心,也积累了经验,为后面的迁移打下了基础。

2. 第二批:中等复杂度任务全面迁移

第一批成功之后,我们开始第二批迁移,也就是大部分中等复杂度的任务,大约有20个。

这些任务,逻辑比第一批复杂一些,有一些聚合、join、窗口函数等操作,而且有一些任务之间有依赖关系。

我们的迁移策略是:

  • 能用SQL的用SQL:对于逻辑比较清晰、能用SQL表达的任务,用Spark SQL重写,开发效率高,可读性好。
  • 复杂逻辑用DataFrame API:对于逻辑比较复杂、SQL不好表达的任务,用DataFrame API写,更灵活,更可控。
  • 公共逻辑抽取成函数:对于多个任务都用到的公共逻辑,抽取成公共函数,在多个任务中复用,减少重复代码。
  • 依赖任务并行迁移:对于有依赖关系的任务,并行迁移,迁移完一个,马上验证,然后迁移下一个。

第二批迁移,遇到的问题比第一批多一些,主要有:

  • 数据倾斜:有些任务,在旧系统中就有数据倾斜的问题,迁移到Spark之后,数据倾斜的问题更明显了,导致某些Task跑得很慢。我们用加盐、两阶段聚合、广播小表等方法,解决了数据倾斜的问题。
  • Shuffle性能问题:有些任务,Shuffle的数据量很大,性能很差。我们通过增加Shuffle分区数、使用Kryo序列化、调整内存比例等方法,优化了Shuffle性能。
  • 内存溢出:有些任务,因为数据量大,或者处理不当,出现了内存溢出。我们通过增加Executor内存、优化数据结构、避免数据倾斜等方法,解决了内存溢出的问题。
  • 结果不一致:有些任务,迁移之后,结果和旧系统有细微的不一致,排查了很久,才发现是空值处理、四舍五入、日期处理等细节问题导致的。我们逐一调整,确保结果完全一致。

第二批任务,花了一个半月的时间,全部迁移完成,数据校验通过,性能也提升了3-5倍。

3. 第三批:核心复杂任务攻坚

第二批完成之后,就剩下第三批了,也就是最复杂、最核心、依赖最多的任务,大约有10个。

这些任务,是整个系统的核心,逻辑非常复杂,有的任务有几千行代码,涉及到多个数据源的join、复杂的业务逻辑、多层的聚合等。而且,这些任务的业务重要性很高,不能出任何问题。

对于这些核心任务,我们采取了更谨慎的策略:

  • 深入理解旧代码:迁移之前,深入阅读旧代码,理解每一行代码的逻辑,确保完全理解之后,再开始迁移。
  • 逐模块迁移:把大任务拆分成小模块,一个模块一个模块地迁移,迁移完一个模块,就验证一个模块,确保每个模块都正确。
  • 双人评审:每个核心任务的代码,都要经过两个人的评审,确保代码质量,没有遗漏。
  • 充分测试:核心任务,要做更充分的测试,包括单元测试、集成测试、数据校验、性能测试等,确保万无一失。
  • 灰度上线:核心任务,先灰度上线,让一小部分数据走新系统,观察有没有问题,没问题再全量上线。

第三批迁移,遇到的问题最多,也最复杂,主要有:

  • 复杂业务逻辑的还原:有些旧代码,逻辑很复杂,而且写得很绕,没有注释,理解起来很困难。我们花了很多时间,逐行阅读旧代码,和业务方沟通,才完全理解了业务逻辑,然后用Spark重写。
  • 多层聚合的性能问题:有些任务,有多层聚合,每层都要Shuffle,性能很差。我们通过优化聚合逻辑、合并聚合、使用窗口函数等方法,减少了Shuffle次数,提升了性能。
  • 大表join的性能问题:有些任务,要join几个大表,数据量很大,性能很差。我们通过广播小表、分桶join、增加资源等方法,优化了join性能。
  • 结果的精确一致性:核心任务,对结果的一致性要求很高,不能有任何差异。我们花了很多时间,做数据校验,逐字段对比,发现了很多细微的差异,比如浮点数精度、日期边界、空值处理等,逐一调整,确保结果完全一致。

第三批任务,花了一个月的时间,全部迁移完成,数据校验通过,性能也提升了3倍左右。

三、数据校验:确保迁移的正确性

大数据迁移,最关键的一点,就是确保新系统的结果和旧系统完全一致。如果结果不一致,迁移就是失败的。所以,数据校验,是迁移过程中非常重要的一环。

我们的数据校验,分为几个层次:

1. 行数校验

最基础的校验,就是行数校验,看看新系统输出的行数,和旧系统是不是一样。如果行数不一样,说明肯定有问题,需要排查。

行数校验,虽然简单,但能发现很多明显的问题,比如过滤条件错了、join的时候丢了数据、重复了数据等。

2. 字段级校验

行数一致之后,还要做字段级校验,看看每个字段的值,是不是一致。我们写了一个数据校验工具,把新旧系统的结果,按主键关联,然后逐字段对比,找出不一致的记录和字段。

字段级校验,能发现很多细微的问题,比如数据类型转换错误、空值处理不一致、四舍五入差异、日期格式差异等。

对于数值型字段,我们还设置了一个容忍度,比如浮点数,允许有0.0001的误差,因为浮点数的精度问题,完全一致可能很难,但误差在容忍度范围内,就是可以接受的。

3. 聚合指标校验

除了明细数据的校验,我们还做了聚合指标的校验,比如总数、总和、平均值、最大值、最小值等,看看新旧系统的聚合指标是不是一致。

聚合指标校验,能从整体上验证数据的正确性,如果聚合指标一致,说明数据的整体分布是对的,即使有个别记录不一致,影响也不大。

4. 业务指标校验

最后,我们还和业务方一起,做了业务指标的校验。比如,某个业务报表的核心指标,新旧系统是不是一致,业务方确认无误之后,才算通过。

业务指标校验,是最关键的一步,因为技术上的校验通过了,不代表业务上就正确了,有些业务逻辑的问题,技术校验可能发现不了,需要业务方来确认。

通过这四个层次的校验,我们确保了迁移之后,新系统的结果和旧系统完全一致,没有任何问题。

四、性能调优:让新系统跑得更快

迁移到Spark,一个重要的目标,就是提升性能。所以,迁移完成之后,我们还做了大量的性能调优工作,让新系统跑得更快。

我们的性能调优,主要从以下几个方面入手:

1. 资源配置调优

首先是资源配置调优,包括Executor的数量、内存、核数,Driver的内存,并行度等。

我们根据每个任务的数据量和复杂度,调整资源配置。数据量大、复杂的任务,给更多的资源;数据量小、简单的任务,给少一些资源,避免浪费。

我们还调整了并行度,也就是Spark的分区数。分区数太少,并行度不够,跑得慢;分区数太多,每个分区的数据量太小,调度开销大。一般来说,分区数设置为Executor总核数的2-3倍,比较合适。

2. 序列化调优

Spark默认的序列化方式是Java序列化,性能比较差,序列化后的体积也比较大。我们改成了Kryo序列化,性能更好,序列化后的体积更小。

使用Kryo序列化,需要注册自定义的类,这样性能更好。我们把项目中用到的自定义类,都注册到了Kryo中。

3. Shuffle调优

Shuffle是Spark性能的瓶颈之一,很多任务慢,都是因为Shuffle慢。我们做了以下Shuffle调优:

  • 增加Shuffle分区数,让每个分区处理的数据量小一些。
  • 调整Shuffle的内存比例,让Shuffle有更多的内存可用。
  • 使用mapSideCombine,在map端先做局部聚合,减少Shuffle的数据量。
  • 对于小表join大表,使用广播join,避免Shuffle。
  • 对于经常join的表,使用分桶表,避免Shuffle。

4. 数据倾斜处理

数据倾斜,是Spark中最常见的性能问题。我们遇到数据倾斜,一般用以下方法处理:

  • 加盐:对倾斜的key加随机前缀,拆分成多个key,分散到不同的Task。
  • 两阶段聚合:先局部聚合,再全局聚合。
  • 广播小表:如果是join导致的倾斜,把小表广播出去。
  • 过滤异常key:如果某些key是异常数据,直接过滤掉。
  • 拆分倾斜key:把倾斜的key单独拿出来处理,然后和其他key的结果合并。

5. 代码优化

除了参数调优,我们还做了代码层面的优化:

  • 避免不必要的Shuffle:尽量减少groupBy、join、distinct等会导致Shuffle的操作。
  • 避免重复计算:对多次使用的RDD或者DataFrame,进行缓存,避免重复计算。
  • 尽早过滤:在计算的早期,就过滤掉不需要的数据,减少后续处理的数据量。
  • 避免在Driver端收集大量数据:不要用collect()把大量数据收集到Driver端,会导致Driver内存溢出。
  • 使用DataFrame API代替RDD API:DataFrame API有 Catalyst优化器,性能比RDD API好很多。

通过这些性能调优,我们的任务,平均性能提升了3-5倍,有的任务甚至提升了10倍以上。

五、遇到的问题和解决方案

在整个迁移过程中,我们遇到了各种各样的问题,这里总结几个比较典型的问题和解决方案。

问题1:新旧系统结果不一致,排查了很久

有一个任务,迁移之后,结果和旧系统不一致,有几千条记录对不上。我们排查了很久,从数据来源、过滤条件、join逻辑、聚合逻辑,逐一排查,都没发现问题。

最后,才发现是日期处理的问题。旧系统用的是MapReduce,自己写的日期处理逻辑,把字符串转成日期,用的是本地时区;而Spark SQL默认用的是UTC时区,导致有一些跨天的记录,日期不一样,所以结果不一致。

解决方案:在Spark SQL中,指定时区为本地时区,或者在处理日期的时候,明确指定时区,确保和旧系统一致。

问题2:内存溢出,加了内存还是溢出

有一个任务,经常出现内存溢出,我们把Executor内存从8G加到16G,又加到32G,还是溢出。

后来排查发现,是因为数据倾斜,某个key对应的数据量特别大,导致这个Task处理的数据量太大,内存不够用。加内存只是治标不治本,因为数据量还在增长。

解决方案:用加盐的方法,把倾斜的key拆分成多个key,分散到不同的Task,每个Task处理的数据量就小了,不会再内存溢出。

问题3:任务运行时间不稳定,有时候快有时候慢

有一个任务,运行时间很不稳定,有时候10分钟跑完,有时候要30分钟,甚至更久。

排查发现,是因为集群资源不稳定,有时候集群负载高,资源不够,任务就跑得慢;有时候集群负载低,资源充足,任务就跑得快。

解决方案:给这个任务配置了资源队列,保证有足够的资源;同时,调整了任务的调度时间,避开集群负载高峰。这样,任务的运行时间就稳定了。

问题4:小文件太多,影响性能

迁移之后,发现新系统输出的小文件特别多,每个任务输出几千个小文件,导致后续读取这些文件的时候,性能很差,NameNode的压力也很大。

原因是Spark的分区数太多,每个分区输出一个文件,所以小文件很多。

解决方案:在输出之前,做一次coalesce或者repartition,减少分区数,从而减少输出文件的数量。同时,我们还调整了表的分区策略,按日期分区,每个分区的文件数量控制在合理范围内。

六、写在最后

这次大数据迁移,从旧的MapReduce系统,迁移到新的Spark系统,花了三个月的时间,迁移了几十个任务,性能提升了3-5倍,开发效率也大大提升。整个过程,踩了很多坑,也积累了很多经验。

总结一下,大数据迁移的几个关键点:

  1. 迁移前的准备很重要:梳理旧系统全貌,评估工作量和优先级,搭建新系统环境,制定迁移规范。
  2. 分批次迁移,稳扎稳打:先迁移简单任务练手,再迁移中等复杂度任务,最后迁移核心复杂任务,风险可控。
  3. 数据校验是关键:确保新系统的结果和旧系统完全一致,从行数、字段、聚合指标、业务指标,多层次校验。
  4. 性能调优不能少:从资源配置、序列化、Shuffle、数据倾斜、代码优化等方面,提升性能。
  5. 遇到问题不要慌:大数据迁移,遇到问题是正常的,耐心排查,找到根本原因,针对性解决。

大数据迁移,是一个复杂的工程,但只要准备充分,方法得当,稳扎稳打,就一定能成功。

希望这次迁移的实战经验,能对正在做大数据迁移的朋友有所帮助。如果有什么问题或者不同的看法,欢迎在评论区交流。