PySpark DataFrame 入門

使用 PySpark 的 Big Data 基礎

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 的 Big Data 基礎

SparkSession:DataFrame API 的入口

  • SparkContext 是建立 RDD 的主要進入點

  • SparkSession 提供與 Spark DataFrame 互動的單一入口

  • SparkSession 用來建立 DataFrame、註冊 DataFrame、執行 SQL 查詢

  • 在 PySpark shell 中,SparkSession 名為 spark

使用 PySpark 的 Big Data 基礎

在 PySpark 建立 DataFrame

  • 在 PySpark 建立 DataFrame 的兩種方法

    • 由現有 RDD,使用 SparkSession 的 createDataFrame() 方法

    • 由各種資料來源(CSV、JSON、TXT),使用 SparkSession 的 read 方法

  • Schema 控制資料並協助 DataFrame 最佳化查詢

  • Schema 提供欄位名稱、欄位資料型別、是否可為空等資訊

使用 PySpark 的 Big Data 基礎

從 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 的 Big Data 基礎

從 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=TrueinferSchema=True
使用 PySpark 的 Big Data 基礎

一起來練習吧!

使用 PySpark 的 Big Data 基礎

Preparing Video For Download...