提升导入性能

使用 PySpark 进行数据清洗

Mike Metzger

Data Engineering Consultant

Spark 集群

Spark 集群由两类进程组成

  • Driver 进程
  • Worker 进程
使用 PySpark 进行数据清洗

导入性能

关键参数:

  • 对象数量(文件、网络位置等)
    • 多个小对象优于单个大对象
    • 可用通配符导入
      airport_df = spark.read.csv('airports-*.txt.gz')
      
  • 对象的大致大小
    • 对象大小相近时 Spark 表现更好
使用 PySpark 进行数据清洗

Schema

明确定义的 schema 可显著提升导入性能

  • 避免多次读取数据
  • 导入时提供校验
使用 PySpark 进行数据清洗

如何拆分对象

  • 使用操作系统工具/脚本(split、cut、awk)
    split -l 10000 -d largefile chunk-
    
  • 使用自定义脚本
  • 写出为 Parquet
    df_csv = spark.read.csv('singlelargefile.csv')
    df_csv.write.parquet('data.parquet')
    df = spark.read.parquet('data.parquet')
    
使用 PySpark 进行数据清洗

Passons à la pratique !

使用 PySpark 进行数据清洗

Preparing Video For Download...