Data Engineering Zoomcamp Module 6 批处理作业实战:用 PySpark 分析 2025 年 11 月纽约出租车数据

Data Engineering Zoomcamp Module 6 批处理作业实战:用 PySpark 分析 2025 年 11 月纽约出租车数据 Data Engineering Zoomcamp Module 6 批处理作业实战用 PySpark 分析 2025 年 11 月纽约出租车数据【免费下载链接】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数据工程训练营Module 6批处理Batch Processing的课后作业完整实战指南。作业要求基于官方提供的 Yellow Taxi 2025-11 Parquet 数据集完成从 PySpark 环境搭建、读取数据、重分区写出到基于 DataFrame/SQL 的过滤、聚合、连接等真实批处理分析任务。读完本文你将掌握 Spark 本地会话的创建方式、Parquet 文件大小估算、日期过滤与时间差计算、Spark UI 端口定位以及通过 zone lookup 表做 join 分组统计的完整实操方法。作业背景与数据集准备Module 6 聚焦批处理技术栈核心是 Apache Spark / PySpark。这份 2026 届作业对应仓库 cohorts/2026/06-batch/homework.md要求把课程中学到的 Spark 知识落实到真实数据上使用的数据集是官方 NYC TLC 发布的Yellow Taxi Trip Data 2025 年 11 月数据Parquet 格式下载命令wget https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2025-11.parquet同时还需要下载出租车区域zone查找表用于第 6 题的区域 join 分析wget https://d37ci6vzurychx.cloudfront.net/misc/taxi_zone_lookup.csv整个作业围绕 6 个问题展开覆盖了 Spark 日常批处理中最常用的一批能力会话创建、数据读取、重分区、过滤、聚合、时间函数与连接查询。环境准备安装 Spark 与 PySpark作业第 1 题是环境搭建安装 Spark、启动 PySpark、创建本地 Spark 会话并执行spark.version输出即为该题的答案具体版本号取决于你安装的版本。依赖与前置条件从仓库的安装指南06-batch/setup/linux.md、06-batch/setup/macos.md、06-batch/setup/windows.md可以确认当前课程使用Spark 4.x它要求Java 17 或 21Linuxsudo apt update sudo apt install default-jdk官方在 Ubuntu 24.04/WSL 上测试通过macOSbrew install openjdk17Windows下载 Adoptium Temurin JDK 17 并解压到工具目录安装完成后建议配置JAVA_HOME并加入PATH再用java --version验证。以 Linux 为例export JAVA_HOME$(dirname $(dirname $(readlink -f $(which java)))) export PATH${JAVA_HOME}/bin:${PATH}安装 PySparkuv 或 pip课程推荐用 uv 管理 Python 依赖uv init uv add pyspark uv run python your_script.py也可以直接用 pippip install pyspark需要特别注意这两种方式安装的 PySpark 都会自带一份打包好的 Spark 发行版无需再单独下载 Spark。如果以前装过 Spark 3.x 并设置了SPARK_HOME环境变量必须把它从.bashrc/.zshrc中移除——否则 PySpark 4.x 会加载到旧 JAR 导致启动失败。Windows 首次运行可能弹出防火墙提示选择允许即可macOS/Windows 上如出现WARNING: Using incubator modules: jdk.incubator.vector警告可安全忽略。验证安装创建本地 Spark 会话仓库 06-batch/setup/linux.md 提供了一个最小验证脚本这也正是作业第 1 题的核心import pyspark from pyspark.sql import SparkSession spark SparkSession.builder \ .master(local[*]) \ .appName(test) \ .getOrCreate() print(fSpark version: {spark.version}) df spark.range(10) df.show() spark.stop()关键点说明master(local[*])表示以本地模式运行*表示使用本机所有可用 CPU 核心适合单机学习与作业练习appName是应用名称会显示在 Spark UI 上便于区分多个作业spark.version返回当前 Spark 版本号——官方解答cohorts/2026/06-batch/solutions.md给出的示例答案是在其环境中输出的4.1.1你的实际输出以本地安装版本为准课程仓库 06-batch/code/04_pyspark.ipynb 中可以看到同样的会话创建写法是贯穿整个 Module 6 的标准启动方式。如果本机环境搭建遇到困难06-batch/README.md 还提供了在 Google Colab 中运行 Spark 的备选方案但课程建议优先把本地环境配置好。读取 Yellow Taxi 数据并重分区写出 Parquet第 2 题要求把 2025 年 11 月的 Yellow 数据读入 Spark DataFrame重分区为 4 个分区后写回 Parquet然后估算生成的每个.parquet文件平均大小单位 MB。官方解答给出的完整实现df spark.read.parquet(yellow_tripdata_2025-11.parquet) df df.repartition(4) df.write.parquet(yellow_tripdata_2025-11_partitioned)随后用 Python 统计每个输出文件的大小import os, glob parquet_files glob.glob(yellow_tripdata_2025-11_partitioned/*.parquet) for f in parquet_files: size_mb os.path.getsize(f) / (1024 * 1024) print(f{os.path.basename(f)}: {size_mb:.1f} MB)原理解读repartition 与输出文件数这里有几个值得深入理解的点Parquet 是一种列式存储格式读取时 Spark 能按列裁剪column pruning配合谓词下推predicate pushdown可显著减少 IO这也是课程选它作为批处理标准格式的原因repartition(4)会对全量数据执行一次shuffle将数据均匀重分布到 4 个分区保证写出的文件数量与分区数一致每个分区独立写一个 Parquet 文件外加可能的_SUCCESS标志文件因此最终目录下会出现 4 个.parquet文件每个约 24.4 MB——这正是该题最接近的答案25MB。注意如果你用coalesce(1)则会合并为单文件。课程示例 06-batch/code/06_spark_sql.py 中就是这样收尾的df_result.coalesce(1) \ .write.parquet(output, modeoverwrite)作业故意考察repartition就是为了让你体验 shuffle 带来的分区数与输出文件数的对应关系。按日期过滤统计记录数第 3 题统计11 月 15 日只考虑 pickup 时间在当天的行程共有多少次出租车出行。官方解答from pyspark.sql import functions as F count_15 df.filter( F.to_date(F.col(tpep_pickup_datetime)) 2025-11-15 ).count() print(count_15) # 162604答案是162,604次。关键 API 讲解F.to_date(...)把时间戳字符串/列转换为日期类型yyyy-MM-dd用于精确匹配某一天如果直接用时间戳做字符串比较会因为带有时分秒而匹配不上F.col(tpep_pickup_datetime)引用列Yellow 数据中 pickup 时间字段名为tpep_pickup_datetimeGreen 数据为lpep_pickup_datetime课程示例 06-batch/code/06_spark_sql.py 中可以看到两者被统一重命名为pickup_datetime的做法.count()是 action 操作会触发一次完整的 job 执行并返回行数。这个写法等价于 SQL 中的WHERE DATE(tpep_pickup_datetime) 2025-11-15Spark 会将其下推为分区/文件级别的过滤条件提高扫描效率。计算最长行程时长第 4 题计算数据集中最长一次行程的时长小时。核心思路是用 dropoff 时间减去 pickup 时间再换算成小时。官方解答df_with_duration df.withColumn( duration_hours, (F.unix_timestamp(F.col(tpep_dropoff_datetime)) - F.unix_timestamp(F.col(tpep_pickup_datetime))) / 3600 ) longest df_with_duration.agg(F.max(duration_hours)).collect()[0][0] print(f{longest:.1f}) # 90.6答案是90.6 小时。关键 API 讲解F.unix_timestamp(col)将时间戳转换为 Unix 秒数两者相减得到行程总秒数除以 3600 得到小时数F.max(duration_hours)在整列上取最大值agg返回聚合结果collect()[0][0]取出第一行第一列的值这种先造新列、再聚合的模式是 Spark DataFrame 编程中最常见的组合套路与课程中groupby/agg的教学内容一脉相承。补充说明unix_timestamp要求列本身是合法的时间类型或可解析的字符串读取 Parquet 时 Spark 会按 schema 自动推断时间列类型因此这里可以直接计算。Spark UI 与默认端口第 5 题是一个经验性知识点Spark 的 Web UI展示应用运行状态、作业、Stage、Executor 等信息的仪表盘默认运行在哪个本地端口正确答案是4040对应选项中的 4040。这是 Spark 的约定俗成Driver 进程启动后默认在 4040 端口暴露 UI如果 4040 被占用会自动顺延到 4041、4042……在本地模式下启动会话后打开浏览器访问http://localhost:4040即可实时观察作业执行、查看 SQL 执行计划和各 Stage 耗时——这也是课程在讲解 Spark 集群解剖Anatomy of a Spark Cluster时反复强调的排障入口。Zone Lookup 连接与最少上车点分析第 6 题是最有综合性的实战题把taxi_zone_lookup.csv载入 Spark 临时视图结合 Yellow 2025-11 数据找出上车次数最少的 pickup location zone 名称。完整实现zones spark.read.option(header, true).csv(taxi_zone_lookup.csv) zones.createOrReplaceTempView(zones) df.createOrReplaceTempView(trips) spark.sql( SELECT z.Zone, COUNT(*) as cnt FROM trips t JOIN zones z ON t.PULocationID z.LocationID GROUP BY z.Zone ORDER BY cnt ASC LIMIT 5 ).show(truncateFalse)输出最少的前 5 个区域------------------------------------------------ |Zone |cnt| ------------------------------------------------ |Governors Island/Ellis Island/Liberty Island|1 | |Eltingville/Annadale/Princes Bay |1 | |Arden Heights |1 | |Port Richmond |3 | |Rikers Island |4 | ------------------------------------------------答案Governors Island/Ellis Island/Liberty Island尽管有 3 个区域同为 1 次但在给出的候选项中该区域是次数最少的。关键点拆解读取 CSV 时用option(header, true)把首行作为列名否则第一行会被当成数据这也是 06-batch/code/04_pyspark.ipynb 中读取 CSV 的标准写法如需自动推断类型可加inferSchema选项临时视图createOrReplaceTempView将 DataFrame 注册为可被 SQL 引用的临时表同一 SparkSession 内有效spark.sql(...)即可直接用 SQL 书写连接、分组、排序join 键行程表的PULocationIDpickup 区域 ID与 zone 表的LocationID关联这是典型的维度表连接fact dimension模式与 Module 6 课程中 join 章节06-batch/README.md讲授的内容一致聚合与排序GROUP BY z.Zone统计每个区域的上车次数ORDER BY cnt ASC升序排列后取前 5即可一眼看出最少区域。题目的候选选项Governors Island/Ellis Island/Liberty Island、Arden Heights、Rikers Island、Jamaica Bay中Jamaica Bay 并未出现在最低频名单里。提交作业与学习分享完成 6 道题后通过课程官方表单提交答案提交地址与截止时间见作业文档 cohorts/2026/06-batch/homework.md 中的链接。官方解答全文位于 cohorts/2026/06-batch/solutions.md可在独立完成后再对照自查。作业文档还鼓励大家Learning in Public公开学习把学习过程整理成 LinkedIn 或 XTwitter帖子记录你掌握了哪些 Spark 技能PySpark 会话、Parquet 处理、repartition 优化、DataFrame 分析、Spark UI 监控等。这类公开输出既是沉淀知识的方式也是建立个人技术影响力的低成本路径。小结本作业覆盖的 Spark 核心能力题号考察主题核心 API/知识点参考答案Q1环境搭建SparkSession、master(local[*])、spark.version按本地安装版本示例 4.1.1Q2分区与写出repartition(4)、write.parquet、文件大小统计25MBQ3日期过滤F.to_date、filter、count162,604Q4时间差聚合F.unix_timestamp、withColumn、F.max90.6 小时Q5Spark UI默认端口 40404040Q6join 与分组CSV 读取、临时视图、spark.sql、GROUP BY ... ORDER BYGovernors Island/Ellis Island/Liberty Island如果你想进一步深入仓库还提供了丰富的配套资源批量下载全年数据的脚本 06-batch/code/download_data.sh按 taxi type 与年份循环wget下载完整的 Spark SQL 聚合示例 06-batch/code/06_spark_sql.py含 revenue 分组、date_trunc按月聚合、union 多数据源等进阶写法以及连接 GCS/BigQuery 的云端配置示例 06-batch/setup/config/spark-defaults.conf。把这些和本作业的 6 道题结合起来你就能系统性地完成从会跑 Spark到能用 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),仅供参考