Операции со столбцами DataFrame

Очистка данных с помощью PySpark

Mike Metzger

Data Engineering Consultant

Повторение: DataFrame

DataFrame:

  • Состоят из строк и столбцов
  • Неизменяемы
  • Для изменения данных используются трансформации
# 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

Основные трансформации 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

Строковые преобразования столбцов

  • Входят в состав pyspark.sql.functions
    import pyspark.sql.functions as F
    
  • Применяются к столбцу как трансформация
    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()

Вспомогательные функции и трансформации для работы с ArrayType()

.size(<column>) — возвращает длину столбца типа arrayType()

.getItem(<index>) — возвращает элемент столбца-списка по указанному индексу.

Очистка данных с помощью PySpark

Давайте потренируемся!

Очистка данных с помощью PySpark

Preparing Video For Download...