Introducere în PySpark DataFrames

Fundamentele Big Data cu PySpark

Upendra Devisetty

Science Analyst, CyVerse

Ce sunt PySpark DataFrames?

  • PySpark SQL este o bibliotecă Spark pentru date structurate. Oferă informații suplimentare despre structura datelor și calcule

  • Un DataFrame PySpark este o colecție distribuită și imuabilă de date cu coloane denumite

  • Conceput pentru procesarea datelor structurate (ex. baze de date relaționale) și semi-structurate (ex. JSON)

  • API-ul DataFrame este disponibil în Python, R, Scala și Java

  • DataFrames în PySpark acceptă atât interogări SQL (SELECT * from table), cât și metode expresive (df.select())

Fundamentele Big Data cu PySpark

SparkSession - Punct de intrare pentru API-ul DataFrame

  • SparkContext este punctul principal de intrare pentru crearea RDD-urilor

  • SparkSession oferă un punct unic de acces pentru interacțiunea cu Spark DataFrames

  • SparkSession este utilizat pentru a crea DataFrames, a le înregistra și a executa interogări SQL

  • SparkSession este disponibil în shell-ul PySpark ca spark

Fundamentele Big Data cu PySpark

Crearea DataFrames în PySpark

  • Două metode de creare a DataFrames în PySpark

    • Din RDD-uri existente, folosind metoda createDataFrame() a SparkSession

    • Din diverse surse de date (CSV, JSON, TXT) folosind metoda read a SparkSession

  • Schema controlează datele și ajută DataFrames să optimizeze interogările

  • Schema furnizează informații despre numele coloanei, tipul datelor, valorile lipsă etc.

Fundamentele Big Data cu PySpark

Crearea unui DataFrame din RDD

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
Fundamentele Big Data cu PySpark

Crearea unui DataFrame din CSV/JSON/TXT

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")
  • Calea către fișier și doi parametri opționali

  • Doi parametri opționali

    • header=True, inferSchema=True
Fundamentele Big Data cu PySpark

Să exersăm!

Fundamentele Big Data cu PySpark

Preparing Video For Download...