使用 Apache Spark 的数据清洗入门

使用 PySpark 进行数据清洗

Mike Metzger

Data Engineering Consultant

什么是数据清洗?

数据清洗:为数据处理流水线准备原始数据。

数据清洗的可能任务:

  • 重新格式化或替换文本
  • 执行计算
  • 移除无效或不完整数据
使用 PySpark 进行数据清洗

为何用 Spark 进行数据清洗?

典型数据系统的问题:

  • 性能
  • 组织数据流

Spark 的优势:

  • 可扩展
  • 强大的数据处理框架
使用 PySpark 进行数据清洗

数据清洗示例

原始数据:

name age (years) city
Smith, John 37 Dallas
Wilson, A. 59 Chicago
null 215

清洗后:

last name first name age (months) state
Smith John 444 TX
Wilson A. 708 IL
使用 PySpark 进行数据清洗

Spark Schema

  • 定义 DataFrame 的格式
  • 可包含多种数据类型:
    • 字符串、日期、整数、数组
  • 可在导入时过滤脏数据
  • 提升读取性能
使用 PySpark 进行数据清洗

Spark Schema 示例

导入 schema

import pyspark.sql.types
peopleSchema = StructType([
  # Define the name field
  StructField('name', StringType(), True),
  # Add the age field
  StructField('age', IntegerType(), True),
  # Add the city field
  StructField('city', StringType(), True)  
])

读取包含数据的 CSV 文件

people_df = spark.read.format('csv').load(name='rawdata.csv', schema=peopleSchema)
使用 PySpark 进行数据清洗

Passons à la pratique !

使用 PySpark 进行数据清洗

Preparing Video For Download...