PySpark DataFrame 入门

使用 PySpark 的大数据基础

Upendra Devisetty

Science Analyst, CyVerse

什么是 PySpark DataFrame?

  • PySpark SQL 是 Spark 的结构化数据库。它提供更多关于数据与计算结构的信息

  • PySpark DataFrame 是带命名列的不可变分布式数据集合

  • 适用于处理结构化数据(如关系型数据库)和半结构化数据(如 JSON)

  • DataFrame API 支持 Python、R、Scala、Java

  • PySpark 的 DataFrame 既可用 SQL(SELECT * from table),也可用表达式方法(df.select()

使用 PySpark 的大数据基础

SparkSession——DataFrame API 的入口

  • SparkContext 是创建 RDD 的主要入口

  • SparkSession 提供与 Spark DataFrame 交互的统一入口

  • SparkSession 用于创建/注册 DataFrame、执行 SQL 查询

  • 在 PySpark shell 中,SparkSession 名为 spark

使用 PySpark 的大数据基础

在 PySpark 中创建 DataFrame

  • 在 PySpark 中创建 DataFrame 的两种方式

    • 基于现有 RDD,使用 SparkSession 的 createDataFrame()

    • 基于数据源(CSV、JSON、TXT),使用 SparkSession 的 read

  • 模式(Schema)约束数据并帮助优化查询

  • 模式包含列名、数据类型、是否空值等信息

使用 PySpark 的大数据基础

从 RDD 创建 DataFrame

iphones_RDD = sc.parallelize([
    ("XS", 2018, 5.65, 2.79, 6.24),
    ("XR", 2018, 5.94, 2.98, 6.84),
    ("X10", 2017, 5.65, 2.79, 6.13),
    ("8Plus", 2017, 6.23, 3.07, 7.12)
])
names = ['Model', 'Year', 'Height', 'Width', 'Weight']
iphones_df = spark.createDataFrame(iphones_RDD, schema=names)

type(iphones_df)
pyspark.sql.dataframe.DataFrame
使用 PySpark 的大数据基础

通过读取 CSV/JSON/TXT 创建 DataFrame

df_csv = spark.read.csv("people.csv", header=True, inferSchema=True)
df_json = spark.read.json("people.json")
df_txt = spark.read.txt("people.txt")
  • 文件路径与两个可选参数

  • 两个可选参数

    • header=True, inferSchema=True
使用 PySpark 的大数据基础

让我们练习

使用 PySpark 的大数据基础

Preparing Video For Download...