Spark算子的介绍与使用
spark在Python中使用,编写环境变量的代码
import math
import os
import re
from pyspark import SparkContext, SparkConf
if __name__ == '__main__':
os.environ['JAVA_HOME'] = r'D:\app\jdk\jdk1.8.0_181'
os.environ['HADOOP_HOME'] = r'D:\app\hadoop\hadoop3.1.4\hadoop-3.1.4'
os.environ['PYSPARK_PYTHON'] = r'C:\Users\86184\miniconda3\python.exe'
os.environ['PYSPARK_DRIVER_PYTHON'] = r'C:\Users\86184\miniconda3\python.exe'
conf = SparkConf().setMaster("local[*]").setAppName("第一个Spark程序")
sc = SparkContext(conf=conf)
1、常见的转换算子:
# 01-map算子,map 就是把 RDD 里的每一条数据,挨个处理一遍,变成新数据。
# 需求:计算每个元素的立方
list1=[1,2,3,4,5,6,7]
rdd1=sc.parallelize(list1).map(lambda x: x**2)
rdd2=sc.parallelize(list1).map(lambda x:math.pow(x,2))
# rdd1.foreach(print)
# rdd2.foreach(lambda x: print(x))
# 02-flatMap()算子-将二维的数据变成一维的,类似于SQL中explode函数
rdd11=sc.textFile("../../datas/day0516/input/a.txt")
# rdd12=rdd11.flatMap(lambda line: line.split("/")).foreach(print)
# 03-filter算子-过滤数据,留下符合条件的,删掉不符合的
rdd3=sc.textFile("../../datas/day0516/input/b.txt")
# 两种都行
# rdd3.filter(lambda line:re.split(r"\s",line)[2] !='-1' and len(re.split("\\s",line))==4).foreach(print)
# rdd3.filter(lambda line:line.split(" ")[2] !='-1' and len(line.split(" "))==4).foreach(print)
2、触发算子:count、foreach、saveAsTextFile
# 01-count-统计RDD集合中元素的个数,返回一个int值,一般用来统计行数
rdd4=sc.textFile("../../datas/day0516/input/a.txt")
# print(rdd4.count())
# 02-foreach-循环遍历
# 03-saveAsTextFile-保存为文件
rdd31=rdd3.filter(lambda line:re.split(r"\s",line)[2] !='-1' and len(re.split("\\s",line))==4)
# rdd31.saveAsTextFile("../../datas/day0516/output")
3、 其他触发算子:first、take、 collect、reduce
#01-first算子-返回第一个分区的第一个元素
# print(rdd31.getNumPartitions())
# print(rdd31.first())
# 02-take算子-返回RDD集合中的前N个元素【先从第一个分区取,如果不够再从第二个分区取】
# print(rdd31.take(3))
# 03-collect算子-collect () = 把 RDD 的所有数据,收集回来,变成 Python 列表(list)
# print(rdd31.collect())
# 04-reduce算子-将RDD中的每个元素按照给定的聚合函数进行聚合,返回聚合的结果
# rdd11.reduceByKey(lambda x, y: x + y).map(lambda x:x[0]).foreach(print)
# rdd11.flatMap(lambda line: line.split("/")).map(lambda x: (x[0],1)).reduceByKey(lambda x, y: x + y).foreach(print)
# reduceByKey是一个转换算子
# rdd4=sc.parallelize([1,2,3,4,5,6,7,8,9]).reduce(lambda x, y: x + y)
# print(rdd4)
# reduce ():对整个 RDD 全局聚合,最终得到一个结果!
# reduceByKey ():按 key 分组聚合,每个 key 返回一个结果!
# 05-TopN算子:top、takeOrdered
# top算子-求排好序之后的最大的几个值 对RDD中的所有元素降序排序,并返回前N个元素,即返回RDD中最大的前N个元数据
# rdd5=sc.parallelize([1,2,3,4,5,6,7,8,9]).top(4)
# print(rdd5)
# takeOrdered : 求排好序之后的最小的几个值 功能:对RDD中的所有元素升序排序,并返回前N个元素,即返回RDD中最小的前N个元数据
rdd6=sc.parallelize([1,2,3,4,5,6,7,8,9]).takeOrdered(3)
print(rdd6)
sc.stop()
4、其他转换算子
# 01-union算子 把两个 RDD 直接拼接在一起,不合并、不排序、不去重
# join是合并
list1 = [1, 2, 3, 4, 5, 6, 7, 8]
list2 = [5, 6, 7, 8, 9, 10]
# rdd1=sc.parallelize(list1).union(sc.parallelize(list2))
# rdd1.foreach(print)
# 02-distinct算子 去重
rdd2=rdd1.distinct()
rdd2.foreach(print)
# 分组聚合算子:groupByKey、 reduceByKey
rdd3 = sc.parallelize([("a", 1), ("a", 1), ("a", 1), ("b", 1)])
# 03-groupByKey算子 先分组,把所有 value 拉到一起,最后让你自己处理
# rdd3.groupByKey().mapValues(lambda x:sum(x)).foreach(print)
# 04-reduceByKey算子 先局部聚合,再全局聚合,性能碾压 groupByKey
# rdd3.reduceByKey(lambda x,y:x+y).foreach(print)
# 注意:能用reduceByKey就不要用groupByKey+map
# reduceByKey代码更简洁,而且性能会更好
# 排序算子:sortBy、sortByKey
# 05-sortByKey算子 按【key】排序 → 只适用于 KV 键值对 RDD
rdd4 = sc.parallelize([("b", 2), ("a", 1), ("c", 3),("d",4),("e",5)])
# rdd4.sortByKey().foreach(print)
# print(rdd4.sortByKey().collect())
"""
不加collect,是两个分区进行排序
('b', 2)
('c', 3)
('d', 4)
('a', 1)
('e', 5)
加了collect,[('a', 1), ('b', 2), ('c', 3), ('d', 4), ('e', 5)]
"""
# 06-sortBy算子 按【自定义字段】排序 → 什么 RDD 都能用,更灵活
# 普通 RDD
rdd5=sc.parallelize([3,4,1,2])
# print(rdd5.sortBy(lambda x: x).collect()) #[1, 2, 3, 4]
# KV RDD 按 value 排
# print(rdd4.sortBy(lambda x: x[1]).collect())#[('a', 1), ('b', 2), ('c', 3), ('d', 4), ('e', 5)]
# 重分区算子:repartition、coalesce
# 07-repartition算子 可以增多、可以减少分区 → 一定触发 shuffle(洗牌) repartition底层就是 coalesce(shuffle=True)
print(rdd4.repartition(4).getNumPartitions())#4
print(rdd4.repartition(2).getNumPartitions())#2
# 08-coalesce算子
print(rdd4.coalesce(3, shuffle=True).getNumPartitions())#增加分区需要shuffle过程,减少不需要
print(rdd4.coalesce(2).getNumPartitions())
# 算子的其他方面
# 1) 其他KV类型算子
#01.keys:专门给键值对 RDD用,只取出所有 key,丢掉 value
rdd_kv = sc.parallelize([('laoda', 11), ('laoer', 22), ('laosan', 33), ('laosi', 44)], numSlices=2)
rdd1=rdd_kv.keys()
# rdd1.foreach(print)#输出样式:laosan
# 02.values:和 keys 对应,专门获取 KV 键值对 RDD 里的所有 value,丢掉 key
rdd2=rdd_kv.values()
# rdd2.foreach(print)#输出样式:11
# 03.mapValues: 将所有的value拿到之后进行map转换,转换后还是元组,只是元组中的value,进行了变化,,,,只处理 value,key 完全不动
# rdd_kv.mapValues(lambda x:x).foreach(print)#输出这样的:('laosan', 33)
# 04.collectAsMap:把 KV 键值对 RDD 直接转成 Python 字典(dict)
# join / fullOuterJoin / leftOuterJoin / rightOuterJoin
# 2)join方面的算子 (转换算子)
join / fullOuterJoin / leftOuterJoin / rightOuterJoin
#都是按 key 拼接(横向合并),区别在于:匹配不到的数据要不要保留!
# 01.join(内连接 / 交集)只保留两边都能匹配上的数据! 有一边没有 → 直接丢掉!
rdd5 = sc.parallelize([("a", 1), ("b", 2), ("c", 3)])
rdd6 = sc.parallelize([("a", 10), ("b", 20), ("d", 30), ("e", 40)])
print(rdd5.join(rdd6).collect()) #[('b', (2, 20)), ('a', (1, 10))]
# 02.fullOuterJoin 全连接
# print(rdd5.fullOuterJoin(rdd6).collect())#[('b', (2, 20)), ('c', (3, None)), ('a', (1, 10)), ('e', (None, 40)), ('d', (None, 30))]
# 03.leftOuterJoin 左外连接
# print(rdd5.leftOuterJoin(rdd6).collect())#[('b', (2, 20)), ('c', (3, None)), ('a', (1, 10))]
# 04.rightOuterJoin 右外连接
# print(rdd5.rightOuterJoin(rdd6).collect())#[('b', (2, 20)), ('a', (1, 10)), ('e', (None, 40)), ('d', (None, 30))]
# 3)分区算子
# 为什么会有分区算子呢?map和foreach 都是针对 一条一条数据的,不是针对一批数据的,假如我需要针对一个大数据集进行处理,有可能会出问题,比如:
# 01.mapPartitions:mapPartitions 是一整个分区一整个分区处理
"""
性能高
减少函数调用次数,批量处理
适合连接数据库
一个分区只创建一次连接,而不是每条数据一次
大数据量优先用
"""
rdd7=sc.parallelize([1,2,3,4,5,6,7,8,9])
# rdd7.map(lambda x:x+1).foreach(print)
def func(x):
for i in x:
yield i+1
# rdd7.mapPartitions(func).foreach(print)
# 02.foreachParition:
# mapPartitions:处理数据 → 返回新 RDD,做转换
# foreachPartition:执行动作 → 无返回,多用于落地入库、打印统计
def deal_part(iter):
for num in iter:
print(num)
print(rdd7.foreachPartition(deal_part))
更多推荐


所有评论(0)