快速管道
在解析更复杂的数据之前,您的经理希望先看到一个包含基本步骤的简单管道示例。在这个示例中,您需要读取一个数据文件,筛选部分行,新增一个 ID 列,然后将结果写出为 JSON 数据。
spark 上下文已定义,且按惯例将 pyspark.sql.functions 库起别名为 F。
本练习是课程的一部分
使用 PySpark 进行数据清洗
练习说明
- 将文件
2015-departures.csv.gz导入为一个 DataFrame。注意文件已包含表头。 - 仅保留航班时长大于 0 分钟的记录。使用列的索引而不是列名(可通过
.printSchema()查看列名和顺序)。 - 新增一个 ID 列。
- 将结果写出为名为
output.json的 JSON 文档。
交互式实操练习
通过完成这段示例代码来试试这个练习。
# Import the data to a DataFrame
departures_df = spark.____(____, header=____)
# Remove any duration of 0
departures_df = departures_df.____(____)
# Add an ID column
departures_df = departures_df.____('id', ____)
# Write the file out to JSON format
____.write.____(____, mode='overwrite')