什么是 RDD

RDD(Resilient Distributed Dataset)叫做分布式数据集,是Spark 中最基本的数据抽象,它代表一个不可变、可分区、里面的元素可并行计算的集合。RDD具有数据流模型的特点:自动容错、位置感知性调度和可伸缩性。RDD允许用户在执行多个查询时显式地将工作集缓存在内存中,后续的查询能够重用工作集,这极大地提升了查询速度。

简单的来说RDD就是一个集合,一个将集合中数据存储在不同机器上的集合。

一个Partitioner,即RDD的分片函数。

按照“移动数据不如移动计算”的理念,Spark 在进行任务调度的时候,会尽可能地将计算任务分配到其所要处理数据块的存储位置。

集合并行化创建RDD

    SparkContext创建;

    sc = SparkContext("local", "Simple App")

 说明:"local" 是指让Spark程序本地运行,"Simple App" 是指Spark程序的名称,这个名称可以任意(为了直观明了的查看,最好设置有意义的名称)。

    集合并行化创建RDD;

       data = [1,2,3,4]
       rdd = sc.parallelize(data)

    collect算子:在驱动程序中将数据集的所有元素作为数组返回(注意数据集不能过大);

    rdd.collect()

    停止SparkContext。

    sc.stop()

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    #********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口
    sc = SparkContext("local", "Simple App")
    # 2.创建一个1到8的列表List
       data = [1, 2, 3, 4, 5, 6, 7, 8]
       rdd = sc.parallelize(data)
    # 3.通过 SparkContext 并行化创建 rdd


    # 4.使用 rdd.collect() 收集 rdd 的内容。 rdd.collect() 是 Spark Action 算子,在后续内容中将会详细说明,主要作用是:收集 rdd 的数据内容
    rdd.collect()
    # 5.打印 rdd 的内容
    print(result)
    # 6.停止 SparkContext
    sc.stop()
    #********** End **********#

读取外部数据集创建RDD

文本文件RDD可以使用创建SparkContex的textFile方法。此方法需要一个 URI的文件(本地路径的机器上,或一个hdfs://,s3a:// 等 URI),并读取其作为行的集合。

# -*- coding: UTF-8 -*-
from pyspark import SparkContext
 
if __name__ == '__main__':
    #********** Begin **********#
 
    # 1.初始化 SparkContext,该对象是 Spark 程序的入口
    sc=SparkContext("local","Simple App")
 
    # 文本文件 RDD 可以使用创建 SparkContext 的t extFile 方法。此方法需要一个 URI的 文件(本地路径的机器上,或一个hdfs://,s3a://等URI),并读取其作为行的集合
    # 2.读取本地文件,URI为:/root/wordcount.txt
    distFile=sc.textFile("/root/wordcount.txt")
 
    # 3.使用 rdd.collect() 收集 rdd 的内容。 rdd.collect() 是 Spark Action 算子,主要作用是:收集 rdd 的数据内容以数组形式返回
    result=distFile.collect()
 
    # 4.打印 rdd 的内容
    print(result)
 
    # 5.停止 SparkContext
    sc.stop()
 
    #********** End **********#

Transformation - map

map

将原来RDD的每个数据项通过map中的用户自定义函数 f 映射转变为一个新的元素。

map 案例

    sc = SparkContext("local", "Simple App")
    data = [1,2,3,4,5,6]
    rdd = sc.parallelize(data)
    print(rdd.collect())
    rdd_map = rdd.map(lambda x: x * 2)
    print(rdd_map.collect())

说明:rdd1 的元素( 1 , 2 , 3 , 4 , 5 , 6 )经过 map 算子( x -> x*2 )转换成了 rdd2 ( 2 , 4 , 6 , 8 , 10 )。

题目:

需求:使用 map 算子,将rdd的数据 (1, 2, 3, 4, 5) 按照下面的规则进行转换操作,规则如下:

    偶数转换成该数的平方;
    奇数转换成该数的立方。

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    #********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口
    sc = SparkContext("local","App")

    # 2.创建一个1到5的列表List
    data = [_ for _ in range(1,6)]

    # 3.通过 SparkContext 并行化创建 rdd
    rdd = sc.parallelize(data)

    # 4.使用rdd.collect() 收集 rdd 的元素。
    print(rdd.collect())

    """
    使用 map 算子,将 rdd 的数据 (1, 2, 3, 4, 5) 按照下面的规则进行转换操作,规则如下:
    需求:
        偶数转换成该数的平方
        奇数转换成该数的立方
    """
    # 5.使用 map 算子完成以上需求
    li = rdd.map(lambda x : x ** 3 if x % 2 == 1 else x ** 2).collect()

    # 6.使用rdd.collect() 收集完成 map 转换的元素
    print(li)

    # 7.停止 SparkContext
    sc.stop()
    #********** End **********#

Transformation - mapPartitions

map:遍历算子,可以遍历RDD中每一个元素,遍历的单位是每条记录。

mapPartitions:遍历算子,可以改变RDD格式,会提高RDD并行度,遍历单位是Partition,也就是在遍历之前它会将一个Partition的数据加载到内存中。

mapPartitions 案例

def f(iterator):
    list = []
    for x in iterator:
        list.append(x*2)
    return list
if __name__ == "__main__":
    sc = SparkContext("local", "Simple App")
    data = [1,2,3,4,5,6]
    rdd = sc.parallelize(data)
    print(rdd.collect())
    partitions = rdd.mapPartitions(f)
    print(partitions.collect())

输出:

    [1, 2, 3, 4, 5, 6]
    [2, 4, 6, 8, 10, 12]

Transformation - filter

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    #********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口
    sc = SparkContext("local", "Simple App")
    # 2.创建一个1到8的列表List
    data = [1, 2, 3, 4, 5, 6, 7, 8]
    # 3.通过 SparkContext 并行化创建 rdd
    rdd = sc.parallelize(data)
    # 4.使用rdd.collect() 收集 rdd 的元素。
    print(rdd.collect())
    """
    使用 filter 算子,将 rdd 的数据 (1, 2, 3, 4, 5, 6, 7, 8) 按照下面的规则进行转换操作,规则如下:
    需求:
        过滤掉rdd中的奇数
    """
    # 5.使用 filter 算子完成以上需求
    rdd_filter = rdd.filter(lambda x: x%2==0)
    # 6.使用rdd.collect() 收集完成 filter 转换的元素
    print(rdd_filter.collect())
    # 7.停止 SparkContext
    sc.stop()

    #********** End **********#
   

Transformation - flatMap


# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
       #********** Begin **********#
       
    # 1.初始化 SparkContext,该对象是 Spark 程序的入口
    sc = SparkContext('local','app')
    # 2.创建一个[[1, 2, 3], [4, 5, 6], [7, 8, 9]] 的列表List
    data = [[1, 2, 3], [4, 5, 6], [7, 8, 9]]
    # 3.通过 SparkContext 并行化创建 rdd
    rdd = sc.parallelize(data)
    # 4.使用rdd.collect() 收集 rdd 的元素。
    print(rdd.collect())
    """
        使用 flatMap 算子,将 rdd 的数据 ([1, 2, 3], [4, 5, 6], [7, 8, 9]) 按照下面的规则进行转换操作,规则如下:
        需求:
            合并RDD的元素,例如:
                            ([1,2,3],[4,5,6])  -->  (1,2,3,4,5,6)
                            ([2,3],[4,5],[6])  -->  (1,2,3,4,5,6)
        """
    # 5.使用 filter 算子完成以上需求
    li = rdd.flatMap(lambda x : x).collect()
    # 6.使用rdd.collect() 收集完成 filter 转换的元素
    print(li)
    # 7.停止 SparkContext
    sc.stop()
    #********** End **********#

Transformation - distinct

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    #********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口

    sc = SparkContext('local','Simple App')
    # 2.创建一个内容为(1, 2, 3, 4, 5, 6, 5, 4, 3, 2, 1)的列表List

    data = [_ for _ in range(1,7)] + [_ for _ in range(5,0,-1)]
    # 3.通过 SparkContext 并行化创建 rdd
    rdd = sc.parallelize(data)

    # 4.使用rdd.collect() 收集 rdd 的元素
    print(rdd.collect())

    """
       使用 distinct 算子,将 rdd 的数据 (1, 2, 3, 4, 5, 6, 5, 4, 3, 2, 1) 按照下面的规则进行转换操作,规则如下:
       需求:
           元素去重,例如:
                        1,2,3,3,2,1  --> 1,2,3
                        1,1,1,1,     --> 1
       """
    # 5.使用 distinct 算子完成以上需求

    li = rdd.distinct().collect()

    # 6.使用rdd.collect() 收集完成 distinct 转换的元素
    print(li)

    # 7.停止 SparkContext
    sc.stop()

    #********** End **********#

Transformation - sortBy

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    # ********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口

    sc = SparkContext('local','Simple App')
    # 2.创建一个内容为(1, 3, 5, 7, 9, 8, 6, 4, 2)的列表List

    data = [1, 3, 5, 7, 9, 8, 6, 4, 2]

    # 3.通过 SparkContext 并行化创建 rdd

    rdd = sc.parallelize(data)

    # 4.使用rdd.collect() 收集 rdd 的元素

    print(rdd.collect())


    """
       使用 sortBy 算子,将 rdd 的数据 (1, 3, 5, 7, 9, 8, 6, 4, 2) 按照下面的规则进行转换操作,规则如下:
       需求:
           元素排序,例如:
            5,4,3,1,2  --> 1,2,3,4,5
       """
    # 5.使用 sortBy 算子完成以上需求

    li = rdd.sortBy(lambda x : x).collect()

    # 6.使用rdd.collect() 收集完成 sortBy 转换的元素
    print(li)

    # 7.停止 SparkContext
    sc.stop()

    #********** End **********#
   

Transformation - sortByKey

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    # ********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口
    sc = SparkContext('local','Simple App')

    # 2.创建一个内容为[(B',1),('A',2),('C',3)]的列表List

    data = [('B',1),('A',2),('C',3)]

    # 3.通过 SparkContext 并行化创建 rdd

    rdd = sc.parallelize(data)

    # 4.使用rdd.collect() 收集 rdd 的元素
    print(rdd.collect())

    """
       使用 sortByKey 算子,将 rdd 的数据 ('B', 1), ('A', 2), ('C', 3) 按照下面的规则进行转换操作,规则如下:
       需求:
           元素排序,例如:
            [(3,3),(2,2),(1,1)]  -->  [(1,1),(2,2),(3,3)]
       """
    # 5.使用 sortByKey 算子完成以上需求

    li = rdd.sortByKey().collect()

    # 6.使用rdd.collect() 收集完成 sortByKey 转换的元素

    print(li)

    # 7.停止 SparkContext
    sc.stop()

    # ********** End **********#

Transformation - mapValues

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    # ********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口
    sc = SparkContext('local','Simple App')

    # 2.创建一个内容为[("1", 1), ("2", 2), ("3", 3), ("4", 4), ("5", 5)]的列表List

    data = [("1", 1), ("2", 2), ("3", 3), ("4", 4), ("5", 5)]
    # 3.通过 SparkContext 并行化创建 rdd

    rdd = sc.parallelize(data);
    # 4.使用rdd.collect() 收集 rdd 的元素
    print(rdd.collect());

    """
           使用 mapValues 算子,将 rdd 的数据 ("1", 1), ("2", 2), ("3", 3), ("4", 4), ("5", 5) 按照下面的规则进行转换操作,规则如下:
           需求:
               元素(key,value)的value进行以下操作:
                                                偶数转换成该数的平方
                                                奇数转换成该数的立方
    """
    # 5.使用 mapValues 算子完成以上需求

    li = rdd.mapValues(lambda x : x ** 3 if x % 2 else x ** 2).collect()

    # 6.使用rdd.collect() 收集完成 mapValues 转换的元素

    print(li)
    # 7.停止 SparkContext
    sc.stop()

    # ********** End **********#



Transformations - reduceByKey


# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    # ********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口
    sc = SparkContext('local','Simple App')

    # 2.创建一个内容为[("python", 1), ("scala", 2), ("python", 3), ("python", 4), ("java", 5)]的列表List
    data = [("python", 1), ("scala", 2), ("python", 3), ("python", 4), ("java", 5)]

    # 3.通过 SparkContext 并行化创建 rdd

    rdd = sc.parallelize(data)

    # 4.使用rdd.collect() 收集 rdd 的元素
    print(rdd.collect())

    """
          使用 reduceByKey 算子,将 rdd 的数据[("python", 1), ("scala", 2), ("python", 3), ("python", 4), ("java", 5)] 按照下面的规则进行转换操作,规则如下:
          需求:
              元素(key-value)的value累加操作,例如:
                                                (1,1),(1,1),(1,2)  --> (1,4)
                                                (1,1),(1,1),(2,2),(2,2)  --> (1,2),(2,4)
    """
    # 5.使用 reduceByKey 算子完成以上需求

    li = rdd.reduceByKey(lambda x,y : x + y).collect()

    # 6.使用rdd.collect() 收集完成 reduceByKey 转换的元素

    print(li)
    # 7.停止 SparkContext
    sc.stop()

    # ********** End **********#

WordCount - 词频统计

题目:

    对文本文件内的每个单词都统计出其出现的次数。
    按照每个单词出现次数的数量,降序排序。

文本文件内容如下:

hello java
hello python java
hello python python
hello flink
scala scala scala scala scala

(scala,5)
(hello,4)
(python,3)
(java,2)
(flink,1)

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":

    """
        需求:对本地文件系统URI为:/root/wordcount.txt 的内容进行词频统计
    """
    # ********** Begin **********#
    sc = SparkContext("local","app");
    rdd = sc.textFile("/root/wordcount.txt")
    li = rdd.flatMap(lambda x : str(x).split(" ")).map(lambda x : (x,1)).reduceByKey(lambda x,y:x + y).sortBy(lambda x : x[1],False).collect();
    print(li)

    # ********** End **********#

Actions - 常用算子

# -*- coding: UTF-8 -*-
from pyspark import SparkContext

if __name__ == "__main__":
    # ********** Begin **********#

    # 1.初始化 SparkContext,该对象是 Spark 程序的入口

    sc = SparkContext('local','Simple App')
    # 2.创建一个内容为[1, 3, 5, 7, 9, 8, 6, 4, 2]的列表List

    data = [1, 3, 5, 7, 9, 8, 6, 4, 2]

    # 3.通过 SparkContext 并行化创建 rdd
    rdd = sc.parallelize(data)

    # 4.收集rdd的所有元素并print输出
    print(rdd.collect())

    # 5.统计rdd的元素个数并print输出

    print(rdd.count())

    # 6.获取rdd的第一个元素并print输出
    print(rdd.first())

    # 7.获取rdd的前3个元素并print输出
    print(rdd.take(3))

    # 8.聚合rdd的所有元素并print输出
    print(rdd.reduce(lambda x,y : x + y))

    # 9.停止 SparkContext

    sc.stop()
    # ********** End **********#

Logo

「智能机器人开发者大赛」官方平台,致力于为开发者和参赛选手提供赛事技术指导、行业标准解读及团队实战案例解析;聚焦智能机器人开发全栈技术闭环,助力开发者攻克技术瓶颈,促进软硬件集成、场景应用及商业化落地的深度研讨。 加入智能机器人开发者社区iRobot Developer,与全球极客并肩突破技术边界,定义机器人开发的未来范式!

更多推荐