Fundamentele Big Data cu PySpark
Upendra Devisetty
Science Analyst, CyVerse
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())
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
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.
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
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=TrueFundamentele Big Data cu PySpark