DataFrame 列操作

使用 PySpark 进行数据清洗

Mike Metzger

Data Engineering Consultant

DataFrame 速览

DataFrame:

  • 由行和列组成
  • 不可变
  • 通过各种转换操作修改数据
# 返回 name 以 "M" 开头的行 
voter_df.filter(voter_df.name.like('M%'))

# 仅返回 name 和 position 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 进行数据清洗

数据筛选

  • 移除空值
  • 移除异常项
  • 拆分合并来源的数据
  • 使用 ~ 取反
    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...