数据管道简介

使用 PySpark 进行数据清洗

Mike Metzger

Data Engineering Consultant

什么是数据管道?

  • 从源到最终输出的数据处理步骤集合
  • 可包含任意数量的步骤或组件
  • 可跨多个系统
  • 本章聚焦 Spark 内的数据管道
使用 PySpark 进行数据清洗

数据管道长什么样?

  • 输入
    • CSV、JSON、Web 服务、数据库
  • 转换
    • withColumn(), .filter(), .drop()
  • 输出
    • CSV、Parquet、数据库
  • 校验
  • 分析
使用 PySpark 进行数据清洗

管道细节

  • Spark 中无正式定义
  • 通常是完成任务所需的常规 Spark 代码
    schema = StructType([
    StructField('name', StringType(), False),
    StructField('age', StringType(), False)
    ])
    df = spark.read.format('csv').load('datafile').schema(schema)
    df = df.withColumn('id', monotonically_increasing_id())
    ...
    df.write.parquet('outdata.parquet')
    df.write.json('outdata.json')
    
使用 PySpark 进行数据清洗

Passons à la pratique !

使用 PySpark 进行数据清洗

Preparing Video For Download...