Data Engineering Zoomcamp 05-Batch 作业实战:用 PySpark 处理 FHV 出租车数据

发布时间:2026/9/12 15:49:27
Data Engineering Zoomcamp 05-Batch 作业实战:用 PySpark 处理 FHV 出租车数据 Data Engineering Zoomcamp 05-Batch 作业实战用 PySpark 处理 FHV 出租车数据【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp本文以 Data Engineering Zoomcamp 2024 届第 5 周Batch Processing / Spark作业为骨架围绕 FHVFor-Hire Vehicle2019 年 10 月数据集系统讲解 Spark 与 PySpark 的安装验证、显式 Schema 定义、CSV 读取与重分区写 Parquet、基于日期的记录统计、时长聚合、Web UI 端口定位以及维度表 Join 分析六类核心实操技能。读完本文你将能够独立完成该周作业并把同样的方法迁移到任何基于 Spark 的批处理数据管道任务中。说明该作业在 2024 届课程中位于05-batch模块而仓库当前主干目录下对应内容已演进为06-batch本文同时引用两份目录下的资源作为参考。作业背景与数据集本次作业的目标是把课程中学到的 Spark 知识落到实处处理纽约市 FHV电召车辆2019 年 10 月的行程数据。FHV 数据与课程中使用的 Yellow / Green Taxi 数据在列结构上有差异因此作业特意要求像课程中那样为 DataFrame 定义显式 Schema这本身就是本作业最重要的训练点之一。作业涉及两类数据FHV 2019-10 行程数据CSV 压缩包包含hvfhs_license_num、dispatching_base_num、pickup_datetime、dropoff_datetime、PULocationID、DOLocationID、SR_Flag七个字段Zone 区域查找表taxi_zone_lookup.csv包含LocationID、Borough、Zone、service_zone四个字段用于把位置 ID 翻译成可读的区域名称。仓库 download_data.sh 展示了课程中批量下载数据的通用脚本模式它接收TAXI_TYPE与YEAR两个参数按月循环拼接 URL用wget把每个月的 gz 压缩包下载到data/raw/taxi_type/year/month/目录。你可以参照该脚本的思路把数据源换成 FHV 2019-10 文件完成下载。题目一安装 Spark 与 PySpark 并验证版本第一题用于确认环境就绪安装 Spark、运行 PySpark、创建本地 Spark 会话并执行spark.version查看输出版本号。安装步骤仓库 06-batch/setup 目录下提供了三份平台安装指南Linux 安装指南含 WSL以 Ubuntu 24.04 验证Windows 安装指南macOS 安装指南以 Linux 为例核心要点是Spark 依赖 JVM先安装 JavaSpark 4.x 要求 Java 17 或 21再通过包管理工具安装 PySpark。由于pip install pyspark会连带捆绑一份完整的 Spark 发行版因此装好 PySpark 即等于装好 Spark无需单独下载 Spark 二进制包。# 1. 安装 Java sudo apt update sudo apt install default-jdk java --version # 验证例如 OpenJDK 21 # 2. 配置 JAVA_HOME写入 ~/.bashrc 或 ~/.zshrc export JAVA_HOME$(dirname $(dirname $(readlink -f $(which java)))) export PATH${JAVA_HOME}/bin:${PATH} # 3. 安装 PySpark推荐 uv 或 pip 二选一 uv init uv add pyspark # 或者 pip install pyspark创建本地 Spark 会话安装完成后参照 Linux 安装指南 中的测试脚本与 06-batch/code/homework.ipynb 中的会话创建方式验证环境import pyspark from pyspark.sql import SparkSession spark SparkSession.builder \ .master(local[*]) \ .appName(test) \ .getOrCreate() print(fSpark version: {spark.version}) # 输出示例4.1.1取决于实际安装版本 spark.stop()SparkSession.builder是 PySpark 2.0 统一的入口 API.master(local[*])表示以本地模式运行并使用全部可用 CPU 核心.appName(...)为应用命名该名称会显示在 Spark Web UI 中.getOrCreate()则保证同一 JVM 内复用已存在的会话。spark.version输出即为本题答案——不同安装版本输出不同本题答案取决于你实际安装的版本。题目二读取 FHV 数据、重分区并写 Parquet测算平均文件大小本题要求按课程方式用显式 Schema 读取 2019 年 10 月 FHV 数据重分区为 6 个分区后保存为 Parquet并回答生成的.parquet文件平均大小单位 MB最接近哪个选项。定义显式 Schema06-batch/code/homework.ipynb 中定义了课程使用的 FHV Schema——这正是本题要求的像课程中那样的 Schema。使用pyspark.sql.types逐字段声明类型可以让 Spark 在读取时直接按类型解析避免 CSV 推断带来的类型错误from pyspark.sql import types schema types.StructType([ types.StructField(hvfhs_license_num, types.StringType(), True), types.StructField(dispatching_base_num, types.StringType(), True), types.StructField(pickup_datetime, types.TimestampType(), True), types.StructField(dropoff_datetime, types.TimestampType(), True), types.StructField(PULocationID, types.IntegerType(), True), types.StructField(DOLocationID, types.IntegerType(), True), types.StructField(SR_Flag, types.StringType(), True) ])字段含义字段名类型含义hvfhs_license_numStringFHV 牌照编号区分 Uber / Lyft 等平台dispatching_base_numString调度基地编号pickup_datetime/dropoff_datetimeTimestamp上客 / 下客时间PULocationID/DOLocationIDInteger上客 / 下客区域 IDSR_FlagString共享行程标志课程用 FHV 版本为字符串类型注意这里pickup_datetime、dropoff_datetime定义为TimestampType()后续题目三、题目四的时间计算都依赖这一正确类型因此 Schema 定义是整份作业的地基。读取 CSV 并按 6 分区写 Parquetdf spark.read \ .option(header, true) \ .schema(schema) \ .csv(fhv_tripdata_2019-10.csv) # 请替换为实际文件路径 df df.repartition(6) df.write.parquet(data/pq/fhv/2019/10/)关键点解释.option(header, true)跳过 CSV 首行表头.schema(schema)使用上一步显式定义的 Schema读取时不再做类型推断repartition(6)把数据重打散为 6 个分区。写 Parquet 时每个分区对应一个.parquet数据文件因此输出目录中约有 6 个数据文件外加_SUCCESS标记文件write.parquet(...)以列式存储格式落盘Parquet 的列式压缩特性也是 Spark 批处理高性能的核心之一。测算平均文件大小import os, glob parquet_files glob.glob(data/pq/fhv/2019/10/*.parquet) sizes [os.path.getsize(f) / (1024 * 1024) for f in parquet_files] print(f文件数: {len(parquet_files)}) print(f平均大小: {sum(sizes) / len(sizes):.1f} MB)原理推导平均文件大小 ≈ 数据集总大小 / 分区数。分区越多单文件越小分区越少单文件越大。作为方法参考仓库中 2026 届作业题解数据为 Yellow Taxi 2025-11、repartition(4)展示的实测结果是每个文件约 24.4 MB——但请以你自己环境对 FHV 2019-10 的实测结果为准。课程中05_taxi_schema.ipynb06-batch/code/05_taxi_schema.ipynb处理 Green/Yellow 数据时同样使用repartition(4)write.parquet的模式可见重分区控制输出文件数是该课程反复使用的工程惯例。题目三统计 10 月 15 日的行程记录数本题要求统计 10 月 15 日仅考虑当天开始的行程有多少笔出租车行程。解法一DataFrame API沿用 06-batch/code/homework.ipynb 中的经典写法——先用to_date把时间戳截断为日期再过滤等于目标日期后计数from pyspark.sql import functions as F df \ .withColumn(pickup_date, F.to_date(df.pickup_datetime)) \ .filter(pickup_date 2019-10-15) \ .count()F.to_date把pickup_datetime的时间部分归零、只保留日期从而可以精确匹配当天开始的行程若直接拿时间戳与字符串比较则可能因时分秒差异导致漏算或误算。解法二Spark SQL把 DataFrame 注册为临时视图后用 SQL 查询结果一致df.createOrReplaceTempView(fhv_2019_10) spark.sql( SELECT COUNT(1) FROM fhv_2019_10 WHERE to_date(pickup_datetime) 2019-10-15; ).show()[!IMPORTANT] 作业特别提示注意定义 Schema 时的列顺序。CSV 读取时 Spark 按 Schema 中字段的声明顺序与文件表头列对应如果列顺序写错pickup_datetime可能被错误解析成其他类型或字段导致to_date过滤结果失真。这一点同样适用于题目二的 Parquet 写入结果。题目四计算数据集中最长行程的时长小时本题要求找出整个数据集中时长最长的一趟行程并以小时为单位作答。时长计算的正确姿势时间戳在 Spark 内部以 Unix 纪元毫秒存储因此两个时间戳相减得到毫秒差。把它换算成小时需要除以 3600×1000或者先把时间戳 cast 成 Long 再除以 3600from pyspark.sql import functions as F df \ .withColumn(duration_hours, (df.dropoff_datetime.cast(long) - df.pickup_datetime.cast(long)) / 3600) \ .agg(F.max(duration_hours)) \ .collect()[0][0]cast(long)把时间戳转为自 Unix 纪元起的秒数TimestampType底层存毫秒cast到 long 会得到秒级值两列相减即得到行程秒数除以 3600 得到小时数agg(F.max(...))对全表取最大值collect()[0][0]取出标量结果。仓库 cohorts/2026/06-batch/solutions.md 展示了等价的另一写法用F.unix_timestamp(...)直接取 Unix 秒再相减同样除以 3600 得小时。两者原理一致可任选。进阶按天分组取最长行程课程作业的考察点是找出整个数据集最长行程若需按天分解可参考 06-batch/code/homework.ipynb 中 Q4 的扩展写法——先算duration秒再按pickup_date分组取maxdf \ .withColumn(duration, df.dropoff_datetime.cast(long) - df.pickup_datetime.cast(long)) \ .withColumn(pickup_date, F.to_date(df.pickup_datetime)) \ .groupBy(pickup_date) \ .max(duration) \ .orderBy(max(duration), ascendingFalse) \ .limit(5) \ .show()-- 等价 SQL把秒换算成分钟并排序 SELECT to_date(pickup_datetime) AS pickup_date, MAX((CAST(dropoff_datetime AS LONG) - CAST(pickup_datetime AS LONG)) / 60) AS duration_minutes FROM fhv_2019_10 GROUP BY 1 ORDER BY 2 DESC LIMIT 10;题目五Spark Web UI 默认监听端口本题考察对 Spark 运维常识的掌握Spark 的 Web UI应用仪表盘默认运行在本地哪个端口答案4040。这一事实也在仓库 2026 届作业题解 中明确标注Sparks User Interface runs on port 4040 by default。当本地以local[*]模式启动 SparkSession 后浏览器访问http://localhost:4040即可查看当前应用的 Stages、Jobs、Executor 与 SQL 执行计划等运行信息。其余选项80 / 443 / 8080分别是 HTTP、HTTPS 及常见 Web 服务器默认端口与 Spark 无关。值得注意的是同一时刻若已有其他 Spark 应用占用 4040新应用会自动改用 4041、4042……依此类推。题目六找出最不频繁的上客区域本题要求加载 Zone 查找表并注册为临时视图与 FHV 2019-10 数据关联找出上客次数最少的Zone名称。第一步加载 Zone 查找表df_zones spark.read \ .option(header, true) \ .csv(taxi_zone_lookup.csv) df_zones.columns # [LocationID, Borough, Zone, service_zone] df_zones.createOrReplaceTempView(zones)第二步Join 分组统计把行程表按PULocationID与 Zone 表的LocationID关联按Zone分组计数并按升序排列取最小者df.createOrReplaceTempView(fhv_2019_10) spark.sql( SELECT z.Zone, COUNT(1) AS cnt FROM fhv_2019_10 fhv LEFT JOIN zones z ON fhv.PULocationID z.LocationID GROUP BY z.Zone ORDER BY cnt ASC LIMIT 5; ).show(truncateFalse)要点说明用LEFT JOIN而非INNER JOIN避免某些行程的PULocationID在 Zone 表中缺失时整行被丢弃ORDER BY cnt ASC LIMIT 5展示计数最少的几个区域通常会出现并列情况此时再结合题目给出的四个候选选项East Chelsea、Jamaica Bay、Union Sq、Crown Heights North比对实际输出中谁在最少之列若用 DataFrame API等价写法是df.groupBy(PULocationID).count()后再joinZone 表。课程笔记与 06-batch/code/homework.ipynb 中的 Q6 展示了同类位置对分析的进阶形态用CONCAT(pul.Zone, / , dol.Zone)拼接上客/下客区域名统计最频繁的起讫区域对。本题反其道而行——统计最少思路完全同构只是排序方向与分组维度不同。提交方式与验证要点提交表单与截止时间见作业文档 cohorts/2024/05-batch/homework.md 末尾说明以课程网站公布为准官方题解视频链接同样记录在作业文档开头可在提交后对照核对若希望对照参考实现仓库中 2026 届题解 提供了完整可运行代码与答案数据为 Yellow Taxi 2025-11题型结构与 2024 届一致是极佳的自测样例。仓库参考资源索引资源路径用途2024 届作业原文cohorts/2024/05-batch/homework.md本作业题面课程 Spark 章节总览06-batch/README.md视频清单与学习路径Linux 安装指南06-batch/setup/linux.mdJava / PySpark 安装与验证Windows / macOS 安装指南06-batch/setup/windows.md、06-batch/setup/macos.md跨平台安装课程 FHV 作业参考 Notebook06-batch/code/homework.ipynbSchema、读取、统计、Join 全流程出租车 Schema 定义06-batch/code/05_taxi_schema.ipynbGreen/Yellow Schema 参考数据下载脚本06-batch/code/download_data.sh按月批量下载模式Spark SQL 批处理示例06-batch/code/06_spark_sql.py临时视图 SQL 聚合 落盘2026 届作业题解cohorts/2026/06-batch/solutions.md同构题型参考答案完成本作业后你实际上已经掌握了批处理数据管道的全部核心环节环境搭建 → 显式 Schema → 数据读取 → 分区控制与列式存储落盘 → 时间维度的过滤与聚合 → 维度表关联分析。这些能力可以无缝迁移到 06_spark_sql.py 所演示的生产级场景——读取多源数据、统一列名、registerTempTable后用 SQL 做月度收入聚合最后coalesce(1).write.parquet(output, modeoverwrite)输出结果这正是真实数据工程中 Spark 批处理任务的典型形态。【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考