Spark 的见解 & 优化 (四)

本贴最后更新于 2520 天前,其中的信息可能已经物是人非

优化

 这里调优主要针对于开发相关的调优见解,参数相关的调优在次不做赘述。

一:广播频繁使用的变量
在 spark 运算中如果有利用 master 表的数据(只读)且数量较大的场景,那么这种场景利用广播该变量来提升性能。

不使用广播变量的场景:
 那么 spark 就会把该数据以网络传输给各个 task,然后每个 task 都有该数据的副本,task 越多,对该数据副本的网络传输也越多,相对应的会加大网络的损耗。副本过多的话也会加大内存的损耗。
使用广播变量的场景:
 那么 spark 就会把该数据的副本以广播的形式缓存到各个节点上,每个节点的所有 task 执行运算时共享该副本,因为每个节点只驻留一份副本,所有内存开销比不用广播变量的要小许多,而且 spark 使用有效的广播算法来分配广播变量,以降低通信成本。

  SparkConf conf = new SparkConf();
  conf.setMaster("local[1]");
  conf.setAppName("test");
  JavaSparkContext jc = new JavaSparkContext(conf);
  Set masters = new HashSet<>();
  masters.add(1);
  masters.add(2);
  masters.add(3);
  List list = new ArrayList<>();
  list.add(1);
  list.add(2);
  list.add(3);
  list.add(4);
  // 广播变量
  Broadcast<Set<Integer>> broadcastVar = jc.broadcast(masters);
  jc.parallelize(list).filter(x->broadcastVar.value().contains(x)?false:true).foreach(x-> System.out.println(x));
 
  结果:
  4

二:mapPartitions 替代 map,foreachPartition 替代 foreach
 不带 partitions 是对应的函数进行转换/遍历所有数据,带 partitions 是对应的函数进行转换/遍历一个分区的数据,优缺点主要有以下几点:
 1)带 partitions 的增加了并行度,对性能有一定的提升
 2)如果要遍历数据往数据库插入数据,每一条数据都执行(创建连接-> 插入数据-> 释放连接)这样的流程的话,性能是十分低下的,而一个分区的数据只使用一个连接,然后批量插入这个分区的数据,最后再释放连接,这种性能上的提升是显著的

  SparkConf conf = new SparkConf();
  conf.setMaster("local[1]");
  conf.setAppName("test");
  JavaSparkContext jc = new JavaSparkContext(conf);
  List<String> list = new ArrayList<>();
  list.add("sql1");
  list.add("sql2");
  list.add("sql3");
  list.add("sql4");
  jc.parallelize(list).foreachPartition(x->{
	// 获取连接
	Connection connection = getConn();
	x.forEachRemaining(x1->{
        connection.prepareStatement(x1);
	});
	// 释放连接
    connection.close();
  });

三:先使用 map 跟 filter 操作对数据结构进行优化跟缩减
 1)数据源最终都会变成一个广范的且有 schema 的对象,如果是计数/去重相关的操作,如果值是无关紧要的,就缩减为 0/1 等,减少内存占用。
 2)做具体的业务之前先看一下是不是所有字段/属性都是必须的,如果有 10 个字段,你只取其中的 3 个字段的话,就尽量把多余的字段/属性给过滤掉。
 3)上面 2 步在层次嵌套较深且数据量大的运算中有很大提升,因为不但减少了内存占用,Shuffle 传输的性能也有提高

  SparkConf conf = new SparkConf();
  conf.setMaster("local[1]");
  conf.setAppName("test");
  JavaSparkContext jc = new JavaSparkContext(conf);
  List<Map<String,Object>> list = new ArrayList<>();
  Map<String,Object>  item = new HashMap<>();
  item.put("a","a1");
  item.put("b","b1");
  item.put("c","c1");
  item.put("d","d1");
  item.put("e","e1");
  list.add(item);
  item = new HashMap<>();
  item.put("a","a11");
  item.put("b","b11");
  item.put("c","c11");
  item.put("d","d11");
  item.put("e","e11");
  item.put("f","e11");
  list.add(item);
  item = new HashMap<>();
  item.put("b","b21");
  item.put("d","d21");
  item.put("e","e21");
  list.add(item);
  // 统计key=a出现的个数
  jc.parallelize(list)
	// 把map的每一条键值对变成多条记录
	.flatMap(x->x.entrySet().iterator())
	// 过滤a以外的key
	.filter(x->"a".equals(x.getKey()))
	// 因为统计key,value从string->int
	.mapToPair(x->new Tuple2<>(x.getKey(),1))
	// 统计
	.reduceByKey((var1,var2)->var1+var2)
	// 遍历
	.foreachPartition(x-> x.forEachRemaining(x1->System.out.println(x1)));
  
  结果:
  (a,2)

四:使用 filter 后进行 coalesce
 如果对 RDD 使用 filter 算子过渡掉较多的数据后,建议使用 coalesce 算子对分区进行缩减操作,使用 coalesce 缩减分区不会产生 Shuffle 操作,适用于数据量较少但分区数较多的场景

  SparkConf conf = new SparkConf();
  conf.setMaster("local[1]");
  conf.setAppName("test");
  JavaSparkContext jc = new JavaSparkContext(conf);
  List<Map<String,Object>> list = new ArrayList<>();
  Map<String,Object>  item = new HashMap<>();
  item.put("a","a1");
  item.put("b","b1");
  item.put("c","c1");
  item.put("d","d1");
  item.put("e","e1");
  list.add(item);
  item = new HashMap<>();
  item.put("a","a11");
  item.put("b","b11");
  item.put("c","c11");
  item.put("d","d11");
  item.put("e","e11");
  item.put("f","e11");
  list.add(item);
  item = new HashMap<>();
  item.put("b","b21");
  item.put("d","d21");
  item.put("e","e21");
  list.add(item);
  // 统计key为a的出现的个数
  JavaRDD, Object>> filter = jc.parallelize(list, 4)
	// 把map的每一条键值对变成多条记录
	.flatMap(x -> x.entrySet().iterator())
	// 过滤a以外的key
	.filter(x -> "a".equals(x.getKey()));
  System.out.println("coalesce前分区数:"+filter.partitions().size());
  // coalesce
  JavaRDD, Object>> coalesce = filter.coalesce(1);
  System.out.println("coalesce后分区数:"+coalesce.partitions().size());
 
  结果:
  coalesce前分区数:4
  coalesce后分区数:1

五:使用 Kryo 进行序列化/反序列化
 spark2.0.0 之前,默认使用的是 java 自带的序列化机制,同时 spark 也支持高性能的 kryo 序列化机制(官方说明比 java 的性能要快 10 倍),但是如果自定义的对象需要在 kryo 进行提前注册才能获得最佳性能。从 spark2.0.0 开始,spark 已经对 kryo 的 github 在这里,大家有兴趣可以自行研究一下。

  // student类
  public class Student {
   private String name;
   private int age;
   private String sex;

   public Student(String name, int age, String sex) {
     this.name = name;
	 this.age = age;
	 this.sex = sex;
   }

   public String getName() {
        return name;
   }

   public void setName(String name) {
        this.name = name;
   }

   public int getAge() {
        return age;
   }

   public void setAge(int age) {
        this.age = age;
   }

   public String getSex() {
        return sex;
   }

   public void setSex(String sex) {
        this.sex = sex;
   }
  }

  // spark示例代码
  SparkConf conf = new SparkConf();
  conf.setMaster("local[1]");
  conf.setAppName("test");
  conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer");
  // conf.registerKryoClasses(new Class[]{Student.class});
  JavaSparkContext jc = new JavaSparkContext(conf);
  List<Student> list = new ArrayList<>();
  for(int i=0;i<1000000;i++){
	list.add(new Student("name_"+i,i,"male"));
  }
  JavaRDD students = jc.parallelize(list);
  // 持久化
  students.persist(StorageLevel.MEMORY_ONLY_SER());
  System.out.println(students.count());
  
  结果:
  使用kryo序列化策略且没有提前注册自定义类的情况下:
  INFO MemoryStore: Block rdd_0_0 stored as bytes in memory (estimated size 28.5 MB, free 883.8 MB)
  使用kryo序列化策略且提前注册自定义类的情况下:
  INFO MemoryStore: Block rdd_0_0 stored as bytes in memory (estimated size 20.9 MB, free 891.4 MB)

六:对多次使用的 RDD 进行持久化
 spark 多次复用 RDD 的原理:每当你对这个 RDD 进行算子操作的时候,spark 都会从源头把这个 RDD 再算出来,然后再对这个 RDD 再进行算子操作,这种情况下,不对这个复用的 RDD 进行持久化操作的话,性能就很差。如果把这个 RDD 持久化到磁盘或者内存后,再针对于这个 RDD 进行算子操作的时候,直接从磁盘或者内存中提取这个 RDD,因为不用再从源头重新计算,这样效率就高了许多。五的示例已经有了持久化的示例,故不再写相关的示例代码。spark 对于持久化的策略有许多种,详细的请参考 org.apache.spark.storage.StorageLevel 这个类
七:尽量避免 Shuffle 算子
 前一篇已经对 Shffle 操作,以及宽窄依赖做了比较详细的说明,所以尽量用非 Shuffle 的算子 + 广播变量等形式来完成业务,这样就可以尽量避免大规模的 IO 操作以及网络传输,可以大大减少性能开销。
八:使用 map-side 预聚合的 shuffle 操作
 如果 Shuffle 算子是不可避免的,那么尽量使用 map-side 预聚合的算子。
 所谓的 map-side 预聚合,就是每个节点都是对相同的 key 进行一次聚合操作,进行 map-side 预聚合后,每个节点都只会留下一个相同的 key,其他节点在拉取相同的 key 的时候,就会大大缩减拉取的数据,从而达到缩减减小 IO 跟网络开销的目的。
 下面以 groupByKey 跟 reduceByKey 进行某个 key 的 sum 统计举例说明:
 groupByKey:

 reduceByKey:

 groupByKey 示例代码:

  SparkConf conf = new SparkConf();
  conf.setMaster("local[1]");
  conf.setAppName("test");
  JavaSparkContext jc = new JavaSparkContext(conf);
  List<Map<String,Integer>> list = new ArrayList<>();
  Map<String,Integer> map = new HashMap<>();
  map.put("yy",1);
  map.put("jj",1);
  map.put("ww",1);
  list.add(map);
  map = new HashMap<>();
  map.put("yy",1);
  list.add(map);
  map = new HashMap<>();
  map.put("yy",1);
  map.put("jj",1);
  list.add(map);
  map = new HashMap<>();
  map.put("yy",1);
  map.put("jj",2);
  list.add(map);
  map = new HashMap<>();
  map.put("jj",1);
  map.put("ww",2);
  list.add(map);
  map = new HashMap<>();
  map.put("ww",1);
  list.add(map);
  jc.parallelize(list)
	// 遍历所有的key
	.flatMap(x -> x.entrySet().iterator())
	// 组成K,V键值对
	.mapToPair(x -> new Tuple2<>(x.getKey(), x.getValue()))
    // groupByKey
	.groupByKey()
    // sum
	.mapValues(x->sum(x)).foreach(x-> System.out.println(x));

  // sum函数
  private static int sum(Iterable integers){
    int sum = 0;
    if(integers != null){
        Iterator iterator = integers.iterator();
	    while (iterator.hasNext()){
            sum += iterator.next();
		}
    }
    return sum;
  } 

 reduceByKey 示例代码:

  SparkConf conf = new SparkConf();
  conf.setMaster("local[1]");
  conf.setAppName("test");
  JavaSparkContext jc = new JavaSparkContext(conf);
  List<Map<String,Integer>> list = new ArrayList<>();
  Map<String,Integer> map = new HashMap<>();
  map.put("yy",1);
  map.put("jj",1);
  map.put("ww",1);
  list.add(map);
  map = new HashMap<>();
  map.put("yy",1);
  list.add(map);
  map = new HashMap<>();
  map.put("yy",1);
  map.put("jj",1);
  list.add(map);
  map = new HashMap<>();
  map.put("yy",1);
  map.put("jj",2);
  list.add(map);
  map = new HashMap<>();
  map.put("jj",1);
  map.put("ww",2);
  list.add(map);
  map = new HashMap<>();
  map.put("ww",1);
  list.add(map);
  jc.parallelize(list)
	// 遍历所有的key
	.flatMap(x -> x.entrySet().iterator())
	// 组成K,V键值对
	.mapToPair(x -> new Tuple2<>(x.getKey(), x.getValue()))
	// reduceByKey
	.reduceByKey((var1,var2)->var1+var2).foreach(x-> System.out.println(x));

 调优相关的官方资料:http://spark.apache.org/docs/latest/tuning.html

上一篇

  • Spark

    Spark 是 UC Berkeley AMP lab 所开源的类 Hadoop MapReduce 的通用并行框架。Spark 拥有 Hadoop MapReduce 所具有的优点;但不同于 MapReduce 的是 Job 中间输出结果可以保存在内存中,从而不再需要读写 HDFS,因此 Spark 能更好地适用于数据挖掘与机器学习等需要迭代的 MapReduce 的算法。

    74 引用 • 46 回帖 • 563 关注

相关帖子

欢迎来到这里!

我们正在构建一个小众社区,大家在这里相互信任,以平等 • 自由 • 奔放的价值观进行分享交流。最终,希望大家能够找到与自己志同道合的伙伴,共同成长。

注册 关于
请输入回帖内容 ...

推荐标签 标签

  • golang

    Go 语言是 Google 推出的一种全新的编程语言,可以在不损失应用程序性能的情况下降低代码的复杂性。谷歌首席软件工程师罗布派克(Rob Pike)说:我们之所以开发 Go,是因为过去 10 多年间软件开发的难度令人沮丧。Go 是谷歌 2009 发布的第二款编程语言。

    502 引用 • 1397 回帖 • 240 关注
  • PostgreSQL

    PostgreSQL 是一款功能强大的企业级数据库系统,在 BSD 开源许可证下发布。

    23 引用 • 22 回帖
  • 导航

    各种网址链接、内容导航。

    45 引用 • 177 回帖
  • OAuth

    OAuth 协议为用户资源的授权提供了一个安全的、开放而又简易的标准。与以往的授权方式不同之处是 oAuth 的授权不会使第三方触及到用户的帐号信息(如用户名与密码),即第三方无需使用用户的用户名与密码就可以申请获得该用户资源的授权,因此 oAuth 是安全的。oAuth 是 Open Authorization 的简写。

    36 引用 • 103 回帖 • 44 关注
  • Spark

    Spark 是 UC Berkeley AMP lab 所开源的类 Hadoop MapReduce 的通用并行框架。Spark 拥有 Hadoop MapReduce 所具有的优点;但不同于 MapReduce 的是 Job 中间输出结果可以保存在内存中,从而不再需要读写 HDFS,因此 Spark 能更好地适用于数据挖掘与机器学习等需要迭代的 MapReduce 的算法。

    74 引用 • 46 回帖 • 563 关注
  • 工具

    子曰:“工欲善其事,必先利其器。”

    308 引用 • 773 回帖
  • RIP

    愿逝者安息!

    8 引用 • 92 回帖 • 429 关注
  • 深度学习

    深度学习(Deep Learning)是机器学习的分支,是一种试图使用包含复杂结构或由多重非线性变换构成的多个处理层对数据进行高层抽象的算法。

    45 引用 • 44 回帖 • 1 关注
  • Pipe

    Pipe 是一款小而美的开源博客平台。Pipe 有着非常活跃的社区,可将文章作为帖子推送到社区,来自社区的回帖将作为博客评论进行联动(具体细节请浏览 B3log 构思 - 分布式社区网络)。

    这是一种全新的网络社区体验,让热爱记录和分享的你不再感到孤单!

    134 引用 • 1128 回帖 • 93 关注
  • Logseq

    Logseq 是一个隐私优先、开源的知识库工具。

    Logseq is a joyful, open-source outliner that works on top of local plain-text Markdown and Org-mode files. Use it to write, organize and share your thoughts, keep your to-do list, and build your own digital garden.

    8 引用 • 69 回帖 • 6 关注
  • wolai

    我来 wolai:不仅仅是未来的云端笔记!

    2 引用 • 14 回帖 • 6 关注
  • Linux

    Linux 是一套免费使用和自由传播的类 Unix 操作系统,是一个基于 POSIX 和 Unix 的多用户、多任务、支持多线程和多 CPU 的操作系统。它能运行主要的 Unix 工具软件、应用程序和网络协议,并支持 32 位和 64 位硬件。Linux 继承了 Unix 以网络为核心的设计思想,是一个性能稳定的多用户网络操作系统。

    960 引用 • 946 回帖
  • Typecho

    Typecho 是一款博客程序,它在 GPLv2 许可证下发行,基于 PHP 构建,可以运行在各种平台上,支持多种数据库(MySQL、PostgreSQL、SQLite)。

    12 引用 • 67 回帖 • 436 关注
  • 京东

    京东是中国最大的自营式电商企业,2015 年第一季度在中国自营式 B2C 电商市场的占有率为 56.3%。2014 年 5 月,京东在美国纳斯达克证券交易所正式挂牌上市(股票代码:JD),是中国第一个成功赴美上市的大型综合型电商平台,与腾讯、百度等中国互联网巨头共同跻身全球前十大互联网公司排行榜。

    14 引用 • 102 回帖 • 260 关注
  • JRebel

    JRebel 是一款 Java 虚拟机插件,它使得 Java 程序员能在不进行重部署的情况下,即时看到代码的改变对一个应用程序带来的影响。

    26 引用 • 78 回帖 • 693 关注
  • Gzip

    gzip (GNU zip)是 GNU 自由软件的文件压缩程序。我们在 Linux 中经常会用到后缀为 .gz 的文件,它们就是 Gzip 格式的。现今已经成为互联网上使用非常普遍的一种数据压缩格式,或者说一种文件格式。

    9 引用 • 12 回帖 • 203 关注
  • Gitea

    Gitea 是一个开源社区驱动的轻量级代码托管解决方案,后端采用 Go 编写,采用 MIT 许可证。

    5 引用 • 16 回帖 • 3 关注
  • AngularJS

    AngularJS 诞生于 2009 年,由 Misko Hevery 等人创建,后为 Google 所收购。是一款优秀的前端 JS 框架,已经被用于 Google 的多款产品当中。AngularJS 有着诸多特性,最为核心的是:MVC、模块化、自动化双向数据绑定、语义化标签、依赖注入等。2.0 版本后已经改名为 Angular。

    12 引用 • 50 回帖 • 531 关注
  • FreeMarker

    FreeMarker 是一款好用且功能强大的 Java 模版引擎。

    23 引用 • 20 回帖 • 475 关注
  • 阿里云

    阿里云是阿里巴巴集团旗下公司,是全球领先的云计算及人工智能科技公司。提供云服务器、云数据库、云安全等云计算服务,以及大数据、人工智能服务、精准定制基于场景的行业解决方案。

    85 引用 • 324 回帖
  • Vim

    Vim 是类 UNIX 系统文本编辑器 Vi 的加强版本,加入了更多特性来帮助编辑源代码。Vim 的部分增强功能包括文件比较(vimdiff)、语法高亮、全面的帮助系统、本地脚本(Vimscript)和便于选择的可视化模式。

    29 引用 • 66 回帖
  • Log4j

    Log4j 是 Apache 开源的一款使用广泛的 Java 日志组件。

    20 引用 • 18 回帖 • 60 关注
  • ActiveMQ

    ActiveMQ 是 Apache 旗下的一款开源消息总线系统,它完整实现了 JMS 规范,是一个企业级的消息中间件。

    19 引用 • 13 回帖 • 707 关注
  • Thymeleaf

    Thymeleaf 是一款用于渲染 XML/XHTML/HTML5 内容的模板引擎。类似 Velocity、 FreeMarker 等,它也可以轻易的与 Spring 等 Web 框架进行集成作为 Web 应用的模板引擎。与其它模板引擎相比,Thymeleaf 最大的特点是能够直接在浏览器中打开并正确显示模板页面,而不需要启动整个 Web 应用。

    11 引用 • 19 回帖 • 413 关注
  • 星云链

    星云链是一个开源公链,业内简单的将其称为区块链上的谷歌。其实它不仅仅是区块链搜索引擎,一个公链的所有功能,它基本都有,比如你可以用它来开发部署你的去中心化的 APP,你可以在上面编写智能合约,发送交易等等。3 分钟快速接入星云链 (NAS) 测试网

    3 引用 • 16 回帖
  • 支付宝

    支付宝是全球领先的独立第三方支付平台,致力于为广大用户提供安全快速的电子支付/网上支付/安全支付/手机支付体验,及转账收款/水电煤缴费/信用卡还款/AA 收款等生活服务应用。

    29 引用 • 347 回帖 • 2 关注
  • V2Ray
    1 引用 • 15 回帖 • 4 关注