Spark实现二次排序

简介: Spark实现二次排序

1. 实验室名称:

大数据实验教学系统

2. 实验项目名称:

练习 Spark实现二次排序

3. 实验学时:

4. 实验原理:

二次排序就是对于<key,value>类型的数据,不但按key排序,而且每个Key对应的value也是有序的。

 Spark 二次排序解决方案:

 我们只需要将年和月组合起来构成一个Key,将第三列作为value,并使用 groupByKey 函数将同一个Key的所有Value全部分组到一起,然后对同一个Key的所有Value进行排序即可。


5. 实验目的:

掌握Spark RDD实现二次排序的方法。

 掌握groupBy、orderBy等算子的使用。


6. 实验内容:

使用Spark RDD实现二次排序。

  假设我们有以下输入文件(逗号分割的分别是“年,月,总数”):

1.      2018,5,22
2.      2019,1,24
3.      2018,2,128
4.      2019,3,56
5.      2019,1,3
6.      2019,2,-43
7.      2019,4,5
8.      2019,3,46
9.      2018,2,64
10.     2019,1,4
11.     2019,1,21
12.     2019,2,35
13.     2019,2,0

我们期望的输出结果是:

1.      2018-2  64,128
2.      2018-5  22
3.      2019-1  3,4,21,24
4.      2019-2  -43,0,35
5.      2019-3  46,56
6.      2019-4  5

7. 实验器材(设备、虚拟机名称):

硬件:x86_64 ubuntu 16.04服务器

 软件:JDK 1.8,Spark-2.3.2,Hadoop-2.7.3,zeppelin-0.8.1


8. 实验步骤:

8.1 启动Spark集群

在终端窗口下,输入以下命令,启动Spark集群:


1.  $ cd /opt/spark
2.  $ ./sbin/start-all.sh

然后使用jps命令查看进程,确保Spark的Master进程和Worker进程已经启动。


8.2 启动zeppelin服务器

在终端窗口下,输入以下命令,启动zeppelin服务器:

1.  $ zeppelin-daemon.sh start

然后使用jps命令查看进程,确保zeppelin服务器已经启动。


8.3 创建notebook文档

1、首先启动浏览器,在地址栏中输入以下url地址,连接zeppelin服务器。

http://localhost:9090

 2、如果zeppelin服务器已正确建立连接,则会看到如下的zeppelin notebook首页。


6a5f04338c8c4049954c68726974b269.png

3、点击【Create new note】链接,创建一个新的笔记本,并命名为”rdd_demo”,解释器默认使用”spark”,如下图所示。

39b057a23ad346cd82816817b297819a.png


8.4 二次排序实现

1、加载数据集。在zeppelin中执行如下代码。


1.  val inputPath = "file:///data/dataset/data.txt"
2.  val inputRDD = sc.textFile(inputPath)

将光标放在代码单元中任意位置,然后同时按下【shift+enter】键,执行以上代码。

 2、实现二次排序。在zeppelin中执行如下代码。

1.  // 实现二次排序
2.  val sortedRDD = inputRDD
3.       .map(line => {
4.            val arr = line.split(",")
5.            val key = arr(0) + "-" + arr(1)
6.            val value = arr(2)
7.            (key,value)
8.       })
9.       .groupByKey()
10.      .map(t => (t._1, t._2.toList.sortWith(_.toInt < _.toInt).mkString(",")))
11.      .sortByKey(true)          // true:升序,false:降序

将光标放在代码单元中任意位置,然后同时按下【shift+enter】键,执行以上代码。

 3、结果输出。在zeppelin中执行如下代码。

1.  sortedRDD.collect.foreach(t => println(t._1 + "\t" + t._2))

将光标放在代码单元中任意位置,然后同时按下【shift+enter】键,执行以上代码。输出结果如下所示。


1.  2018-2    64,128
2.  2018-5    22
3.  2019-1    3,4,21,24
4.  2019-2    -43,0,35
5.  2019-3    46,56
6.  2019-4    5

8.5 转换为DataFrame输出

也可以将结果转换为DataFrame输出。在zeppelin中执行如下代码。

1.  // 加载数据集
2.  val inputPath = "file:///data/dataset/data.txt"
3.  val inputRDD = sc.textFile(inputPath)
4.       
5.  // 实现二次排序
6.  val sortedRDD = inputRDD
7.                    .map(line => {
8.                      val arr = line.split(",")
9.                      val key = arr(0) + "-" + arr(1)
10.                     val value = arr(2)
11.                     (key,value)
12.                   })
13.                   .groupByKey()
14.                   .map(t => (t._1, t._2.toList.sortWith(_.toInt < _.toInt).mkString(",")))
15.                   .sortByKey(true)          // true:升序,false:降序
16.      
17. // 转换为DataFrame
18. val sortedDF = sortedRDD.toDF("key", "value")
19.      
20. // 显示
21. sortedDF.show

将光标放在代码单元中任意位置,然后同时按下【shift+enter】键,执行以上代码。输出结果如下所示。


9. 实验结果及分析:

实验结果运行准确,无误


10. 实验结论:

经过本节实验的学习,通过练习 Spark实现二次排序,进一步巩固了我们的Spark基础。


11. 总结及心得体会:

二次排序就是对于<key,value>类型的数据,不但按key排序,而且每个Key对应的value也是有序的。

 Spark 二次排序解决方案:

 我们只需要将年和月组合起来构成一个Key,将第三列作为value,并使用 groupByKey 函数将同一个Key的所有Value全部分组到一起,然后对同一个Key的所有Value进行排序即可。

9056a8ab45194ee585321729041b0937.png

相关文章
|
分布式计算 大数据 Spark
Spark textFile 和排序-2
快速学习 Spark textFile 和排序-2
265 0
Spark textFile 和排序-2
|
分布式计算 大数据 Spark
Spark textFile 和排序-4
快速学习 Spark textFile 和排序-4
244 0
|
分布式计算 大数据 Spark
Spark textFile 和排序-3
快速学习 Spark textFile 和排序-3
239 0
|
分布式计算 Java 大数据
Spark textFile 和排序-1
快速学习 Spark textFile 和排序-1
261 0
|
分布式计算 Spark
Spark多路径输出和二次排序
打开微信扫一扫,关注微信公众号【数据与算法联盟】 转载请注明出处:http://blog.csdn.net/gamer_gyt 博主微博:http://weibo.com/234654758 Github:https://github.com/thinkgamer 在实际应用场景中,我们对于Spark往往有各式各样的需求,比如说想MR中的二次排序,Top N,多路劲输出等。
1842 0
|
分布式计算 搜索推荐 Apache
|
人工智能 分布式计算 大数据
大数据≠大样本:基于Spark的特征降维实战(提升10倍训练效率)
本文探讨了大数据场景下降维的核心问题与解决方案,重点分析了“维度灾难”对模型性能的影响及特征冗余的陷阱。通过数学证明与实际案例,揭示高维空间中样本稀疏性问题,并提出基于Spark的分布式降维技术选型与优化策略。文章详细展示了PCA在亿级用户画像中的应用,包括数据准备、核心实现与效果评估,同时深入探讨了协方差矩阵计算与特征值分解的并行优化方法。此外,还介绍了动态维度调整、非线性特征处理及降维与其他AI技术的协同效应,为生产环境提供了最佳实践指南。最终总结出降维的本质与工程实践原则,展望未来发展方向。
766 0
|
分布式计算 大数据 Apache
ClickHouse与大数据生态集成:Spark & Flink 实战
【10月更文挑战第26天】在当今这个数据爆炸的时代,能够高效地处理和分析海量数据成为了企业和组织提升竞争力的关键。作为一款高性能的列式数据库系统,ClickHouse 在大数据分析领域展现出了卓越的能力。然而,为了充分利用ClickHouse的优势,将其与现有的大数据处理框架(如Apache Spark和Apache Flink)进行集成变得尤为重要。本文将从我个人的角度出发,探讨如何通过这些技术的结合,实现对大规模数据的实时处理和分析。
1318 2
ClickHouse与大数据生态集成:Spark & Flink 实战
|
存储 分布式计算 Hadoop
从“笨重大象”到“敏捷火花”:Hadoop与Spark的大数据技术进化之路
从“笨重大象”到“敏捷火花”:Hadoop与Spark的大数据技术进化之路
856 79