การดำเนินการกับคอลัมน์ใน DataFrame

การทำความสะอาดข้อมูลด้วย PySpark

Mike Metzger

Data Engineering Consultant

ทบทวน DataFrame

DataFrames:

  • ประกอบด้วยแถวและคอลัมน์
  • Immutable
  • ใช้การ transformation ต่าง ๆ เพื่อแก้ไขข้อมูล
# Return rows where name starts with "M" 
voter_df.filter(voter_df.name.like('M%'))

# Return name and position only voters = voter_df.select('name', 'position')
การทำความสะอาดข้อมูลด้วย PySpark

การ transformation ที่ใช้บ่อยใน DataFrame

  • Filter / Where
    voter_df.filter(voter_df.date > '1/1/2019') # or voter_df.where(...)
    
  • Select
    voter_df.select(voter_df.name)
    
  • withColumn
    voter_df.withColumn('year', voter_df.date.year)
    
  • drop
    voter_df.drop('unused_column')
    
การทำความสะอาดข้อมูลด้วย PySpark

การกรองข้อมูล

  • ลบค่า null
  • ลบข้อมูลที่ผิดปกติ
  • แยกข้อมูลจากแหล่งรวม
  • ใช้ ~ เพื่อนิเสธเงื่อนไข
    voter_df.filter(voter_df['name'].isNotNull())
    voter_df.filter(voter_df.date.year > 1800)
    voter_df.where(voter_df['_c0'].contains('VOTE'))
    voter_df.where(~ voter_df._c1.isNull())
    
การทำความสะอาดข้อมูลด้วย PySpark

การแปลง String ในคอลัมน์

  • อยู่ใน pyspark.sql.functions
    import pyspark.sql.functions as F
    
  • ใช้กับคอลัมน์เป็น transformation
    voter_df.withColumn('upper', F.upper('name'))
    
  • สร้างคอลัมน์ตัวกลางได้
    voter_df.withColumn('splits', F.split('name', ' '))
    
  • แปลงชนิดข้อมูลได้
    voter_df.withColumn('year', voter_df['_c4'].cast(IntegerType()))
    
การทำความสะอาดข้อมูลด้วย PySpark

ฟังก์ชันสำหรับคอลัมน์ ArrayType()

ฟังก์ชันและการ transformation สำหรับใช้งานกับ ArrayType()

.size(<column>) - คืนค่าความยาวของคอลัมน์ชนิด arrayType()

.getItem(<index>) - ดึงข้อมูลที่ตำแหน่ง index ที่ระบุในคอลัมน์ list

การทำความสะอาดข้อมูลด้วย PySpark

มาฝึกกันเถอะ!

การทำความสะอาดข้อมูลด้วย PySpark

Preparing Video For Download...