理解 Parquet

使用 PySpark 进行数据清洗

Mike Metzger

Data Engineering Consultant

CSV 的难点

  • 无固定模式
  • 嵌套数据需特殊处理
  • 编码格式受限
使用 PySpark 进行数据清洗

Spark 与 CSV

  • 解析慢
  • 无法按条件过滤(不支持"谓词下推")
  • 中间使用需重新定义模式
使用 PySpark 进行数据清洗

Parquet 格式

  • 列式数据格式
  • 受 Spark 等数据处理框架支持
  • 支持谓词下推
  • 自动存储模式信息
使用 PySpark 进行数据清洗

操作 Parquet

读取 Parquet 文件

df = spark.read.format('parquet').load('filename.parquet')
df = spark.read.parquet('filename.parquet')

写入 Parquet 文件

df.write.format('parquet').save('filename.parquet')
df.write.parquet('filename.parquet')
使用 PySpark 进行数据清洗

Parquet 与 SQL

将 Parquet 作为 SparkSQL 的底层存储

flight_df = spark.read.parquet('flights.parquet')
flight_df.createOrReplaceTempView('flights')
short_flights_df = spark.sql('SELECT * FROM flights WHERE flightduration < 100')
使用 PySpark 进行数据清洗

开始练习!

使用 PySpark 进行数据清洗

Preparing Video For Download...