Operații pe coloane DataFrame

Curățarea datelor cu PySpark

Mike Metzger

Data Engineering Consultant

Recapitulare DataFrame

DataFrame-uri:

  • Formate din rânduri și coloane
  • Imuabile
  • Folosesc operații de transformare pentru a modifica datele
# 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')
Curățarea datelor cu PySpark

Transformări comune ale 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')
    
Curățarea datelor cu PySpark

Filtrarea datelor

  • Eliminarea valorilor null
  • Eliminarea înregistrărilor incorecte
  • Separarea datelor din surse combinate
  • Negare cu ~
    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())
    
Curățarea datelor cu PySpark

Transformări pe coloane de tip șir

  • Incluse în pyspark.sql.functions
    import pyspark.sql.functions as F
    
  • Aplicate per coloană ca transformare
    voter_df.withColumn('upper', F.upper('name'))
    
  • Pot crea coloane intermediare
    voter_df.withColumn('splits', F.split('name', ' '))
    
  • Pot converti la alte tipuri
    voter_df.withColumn('year', voter_df['_c4'].cast(IntegerType()))
    
Curățarea datelor cu PySpark

Funcții pentru coloane ArrayType()

Funcții utilitare pentru lucrul cu ArrayType()

.size(<column>) - returnează lungimea coloanei de tip arrayType()

.getItem(<index>) - preia un element specific la indexul dat dintr-o coloană de tip listă.

Curățarea datelor cu PySpark

Să exersăm!

Curățarea datelor cu PySpark

Preparing Video For Download...