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))

更多推荐