使用 PySpark 进行数据清洗
Mike Metzger
Data Engineering Consultant
读取 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')
将 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 进行数据清洗