PySpark RDD 简介

使用 PySpark 的大数据基础

Upendra Devisetty

Science Analyst, CyVerse

什么是 RDD?

  • RDD = Resilient Distributed Datasets(弹性分布式数据集)

rdd

使用 PySpark 的大数据基础

分解 RDD 概念

  • 弹性分布式数据集

    • 弹性:可容错,能应对失败

    • 分布式:跨多台机器

    • 数据集:分区的数据集合,如数组、表、元组等

使用 PySpark 的大数据基础

创建 RDD:如何做?

  • 并行化现有对象集合

  • 外部数据集:

    • HDFS 中的文件

    • Amazon S3 存储桶中的对象

    • 文本文件中的行

  • 基于现有 RDD 创建

使用 PySpark 的大数据基础

并行化集合(parallelizing)

  • 使用 parallelize() 从 Python 列表创建 RDD
numRDD = sc.parallelize([1,2,3,4])
helloRDD = sc.parallelize("Hello world")
type(helloRDD)
<class 'pyspark.rdd.PipelinedRDD'>
使用 PySpark 的大数据基础

来自外部数据集

  • 使用 textFile() 从外部数据集创建 RDD
fileRDD = sc.textFile("README.md")
type(fileRDD)
<class 'pyspark.rdd.PipelinedRDD'>
使用 PySpark 的大数据基础

理解 PySpark 的分区

  • 分区是对大型分布式数据集的逻辑划分

  • parallelize() 方法

numRDD = sc.parallelize(range(10), minPartitions = 6)
  • textFile() 方法
fileRDD = sc.textFile("README.md", minPartitions = 6)
  • 使用 getNumPartitions() 可获取 RDD 的分区数
使用 PySpark 的大数据基础

让我们来练习

使用 PySpark 的大数据基础

Preparing Video For Download...