Big Data Fundamentals with PySpark
Upendra Devisetty
Science Analyst, CyVerse
PySpark SQL คือไลบรารี Spark สำหรับข้อมูลเชิงโครงสร้าง ให้ข้อมูลเพิ่มเติมเกี่ยวกับโครงสร้างข้อมูลและการคำนวณ
PySpark DataFrame คือชุดข้อมูลแบบกระจายที่เปลี่ยนแปลงไม่ได้ มีคอลัมน์ที่มีชื่อ
ออกแบบมาเพื่อประมวลผลข้อมูลเชิงโครงสร้าง (เช่น ฐานข้อมูลเชิงสัมพันธ์) และกึ่งโครงสร้าง (เช่น JSON)
DataFrame API ใช้งานได้ใน Python, R, Scala และ Java
DataFrames ใน PySpark รองรับทั้ง SQL queries (SELECT * from table) และเมธอดแบบ expression (df.select())
SparkContext คือจุดเริ่มต้นหลักสำหรับการสร้าง RDD
SparkSession เป็นจุดเข้าถึงเดียวสำหรับการทำงานกับ Spark DataFrames
SparkSession ใช้สร้าง DataFrame, ลงทะเบียน DataFrames และรัน SQL queries
SparkSession ใน PySpark shell เรียกใช้งานผ่านตัวแปร spark
มี 2 วิธีในการสร้าง DataFrames ใน PySpark
จาก RDD ที่มีอยู่ โดยใช้เมธอด createDataFrame() ของ SparkSession
จากแหล่งข้อมูลต่าง ๆ (CSV, JSON, TXT) โดยใช้เมธอด read ของ SparkSession
Schema ควบคุมข้อมูลและช่วยให้ DataFrames ปรับ query ได้อย่างมีประสิทธิภาพ
Schema ระบุชื่อคอลัมน์, ประเภทข้อมูล, ค่าว่าง และอื่น ๆ
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")
path ของไฟล์และพารามิเตอร์เสริม 2 ตัว
พารามิเตอร์เสริม 2 ตัว
header=True, inferSchema=TrueBig Data Fundamentals with PySpark