{"metadata":{"kernelspec":{"language":"python","display_name":"Python 3","name":"python3"},"language_info":{"name":"python","version":"3.12.12","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceId":31254,"databundleVersionId":3103714,"sourceType":"competition"}],"dockerImageVersionId":31234,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"### 上海应用技术大学2025—2026学年第 1 学期\n### 《大数据分析技术及应用》期末复习大作业\n\n#### 课程代码: B3104516  学分:   3     考试时间:   N/A   分钟\n\n#### 课程序号:                                                                         \n#### 班级：            学号：                姓名:               \n\n#### 作业提交方式：完成后运行所有单元格，确定输出完整，然后下载.ipynb文件，作为附件上传到学习通答题处。\n---\n\n# H&M销售数据分析\n本笔记本以 H&M 个性化推荐真实交易数据 为案例，系统回顾并综合运用本课程中 PySpark 与 Spark ML 的核心知识，完整呈现从原始数据理解到特征工程建模，再到推荐系统训练与调优的一条端到端数据分析与建模流程。笔记本共分为四个部分，既覆盖了课程的主要技术模块，也强调对 Spark 执行机制与模型设计假设的理解，旨在帮助同学在期末复习阶段形成结构化、可迁移的知识体系。\n\n第一部分为探索性数据分析（EDA）。通过对 H&M 交易数据的读取、字段理解与基础统计分析，帮助同学从业务与数据结构两个层面认识数据规模、分布特征与潜在问题，为后续特征工程和模型选择奠定数据理解基础。\n\n第二部分为特征工程。在多表连接与数据抽样的基础上，构建用于预测用户行为的特征工程 Pipeline，系统练习类别型特征编码、数值型特征处理、文本特征分词与清洗等关键操作，并结合 Spark 的缓存机制与执行模型，理解在大规模数据处理中“为什么这样做”和“什么时候值得这样做”。\n\n第三部分为机器学习建模。基于 Spark ML Pipeline，将前一部分构建的特征统一组织为可复用的建模流程，完成模型训练、参数搜索与交叉验证，重点考察对 Pipeline、CrossValidator 以及评估指标的理解，而不仅是模型 API 的使用。\n\n第四部分为推荐系统。围绕协同过滤思想，使用 ALS 模型对用户–物品–交互数据进行建模，完成从用户–物品数据重构、模型调优到推荐结果生成的完整流程，并结合 Spark 的执行日志与警告信息，深入理解分布式推荐模型在实际运行中的调度与性能特征。\n\n通过这四个部分的综合训练，本笔记本不仅要求同学能够“写出能运行的代码”，更强调理解 Spark 在做什么、模型在学什么、结果代表什么。完成本作业后，应能够独立梳理一套基于 PySpark 的数据分析与建模流程，并具备将其迁移到其他真实业务场景中的能力。\n\n---\n## 作业要求\n\n本次作业不禁止学生在作业过程中使用 AI 工具进行辅助学习，但**明确不鼓励直接复制粘贴 AI 生成的完整代码作为作业答案**。AI 工具应主要用于理解概念、澄清思路和排查错误，**而非替代学生完成作业**。\n\n本次作业为期末复习作业，所涉及的代码主要围绕 Spark 数据读取、DataFrame 操作、特征工程 Pipeline 以及常见机器学习与推荐系统流程，代码体量较小、结构明确，**具备在理解课程内容基础上自行完成的可行性**。同时，期末考试为闭卷形式，考试中将重点考察对代码逻辑、关键参数含义及执行流程的理解，而非对完整代码的记忆。**若在复习阶段仅通过复制粘贴方式完成作业，而未经过独立思考与实际编写过程，往往难以在考试中正确理解题意或组织解题步骤，因此无法达到复习与巩固知识的目的**。\n\n本次作业为期末复习性质，主要用于引导学生回顾课程核心内容。为鼓励理解性完成，本次作业采用低完成度满分机制，完成约 60% 即可获得满分，且只计入《课程设计》的平时成绩。\n\n---\n\n# 第一部分：探索性分析\n探索性分析（Exploratory Data Analysis, EDA） 是本实验的第一步，目的是在正式构建推荐模型之前，对 H&M 交易数据的结构、规模与基本特征进行系统性理解。本阶段将围绕销售时间分布、客户行为特征、商品销售情况等关键维度，对数据进行初步统计，以识别数据中存在的规律、差异与潜在问题。通过探索性分析，一方面可以检验数据质量与字段含义是否符合业务直觉，另一方面也为后续特征工程与推荐算法的选择提供依据，确保推荐系统的建模建立在对数据充分理解的基础之上。","metadata":{}},{"cell_type":"markdown","source":"0. 创建SparkSession","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"1. 读取右侧边栏中的transactions_train.csv数据，创建名为transaction的DataFrame展示数据的schema。预期输出为：\n```\nroot\n |-- t_dat: date (nullable = true)\n |-- customer_id: string (nullable = true)\n |-- article_id: integer (nullable = true)\n |-- price: double (nullable = true)\n |-- sales_channel_id: integer (nullable = true)\n```\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"2. 通过创建数据集sample，创建一个新的数据转换计划，数据转换步骤如下：\n   \n   - （1）由于原数据集很大，Spark将transaction分成了很多个分区（Partition），先使用.coalesce(8)将transanction合并为8个分区。\n   - （2）然后，从中随机抽取一个比例为0.00005的样本，以你自己的学号为随机种子。\n   - （3）使用.cache() 将这个样本强制缓存到内存中，减少后续的读写过程。\n   \n   使用.explain(\"\")查看sample所代表的数据转换计划。参考输出如下:\n\n```\n   == Physical Plan ==\nAdaptiveSparkPlan isFinalPlan=false\n+- InMemoryTableScan [t_dat#17, customer_id#18, article_id#19, price#20, sales_channel_id#21]\n      +- InMemoryRelation [t_dat#17, customer_id#18, article_id#19, price#20, sales_channel_id#21], StorageLevel(disk, memory, deserialized, 1 replicas)\n            +- *(1) Sample 0.0, 5.0E-5, false, 123456\n               +- Coalesce 8\n                  +- FileScan csv [t_dat#17,customer_id#18,article_id#19,price#20,sales_channel_id#21] Batched: false, DataFilters: [], Format: CSV, Location: InMemoryFileIndex(1 paths)[file:/kaggle/input/h-and-m-personalized-fashion-recommendations/transa..., PartitionFilters: [], PushedFilters: [], ReadSchema: struct<t_dat:date,customer_id:string,article_id:int,price:double,sales_channel_id:int>\n\n```\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"3. 请解释为什么说Spark DataFrame代表不可变、延迟计算的数据转换计划？第2题中有真正的计算发生吗？","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":" 4. 查看sample的行数。参考输出如下，注意，具体的数字不会完全一样。\n\n```\n1615\n```\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"5. 已知transaction的行数为31,788,324，sample的行数刚好等于它的0.00005吗？这说明了Spark分布式计算的什么机制？","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"6. 展示sample的前5行数据。参考输出如下，注意，具体的数字不会完全一样。\n\n```\n+----------+--------------------+----------+--------------------+----------------+\n|     t_dat|         customer_id|article_id|               price|sales_channel_id|\n+----------+--------------------+----------+--------------------+----------------+\n|2018-09-20|7bb6c88284850a04c...| 685687003|0.016932203389830508|               2|\n|2018-09-20|d906c5406745b9853...| 708021002|0.016932203389830508|               2|\n|2018-09-20|dc7e7d29a490da345...| 662980003|0.033881355932203386|               1|\n|2018-09-21|005777ba7b5f41487...| 539723006|0.033881355932203386|               2|\n|2018-09-21|f90c5ddca53a4f6d1...| 615367010|0.022016949152542376|               2|\n+----------+--------------------+----------+--------------------+----------------+\nonly showing top 5 rows\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"7. （1）从sample中选取物品ID和价格两列，再添加一列“priceYuan”,使其等于price的7倍，将这三列新建一个DataFrame,命名为article\n   \n   （2）查看article的前5行。\n\n   参考输出如下，数字不会完全一样：\n```\n+----------+--------------------+-------------------+\n|article_id|               price|          priceYuan|\n+----------+--------------------+-------------------+\n| 685687003|0.016932203389830508|0.11852542372881356|\n| 708021002|0.016932203389830508|0.11852542372881356|\n| 662980003|0.033881355932203386|0.23716949152542371|\n| 539723006|0.033881355932203386|0.23716949152542371|\n| 615367010|0.022016949152542376|0.15411864406779663|\n+----------+--------------------+-------------------+\nonly showing top 5 rows\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"8. 请问第7题中，触发计算的是步骤（1）还是步骤（2）？为什么会产生这种情况？","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"9. 新建一个DataFrame，命名为customer，它将sample按照客户ID分组，统计price的总合，并重命名为sum_price列，同时分组统计price的均值，重命名为avg_price。展示customer的前10行。\n    预期输出如下，数字不会完全一样。\n\n```\n+--------------------+--------------------+--------------------+\n|         customer_id|           sum_price|          mean_price|\n+--------------------+--------------------+--------------------+\n|89b6b868175c7883e...| 0.02159322033898305| 0.02159322033898305|\n|d906c5406745b9853...|0.016932203389830508|0.016932203389830508|\n|d78fea49cbc32e296...|0.040661016949152536|0.040661016949152536|\n|bb75571162a6f4e26...|0.042355932203389825|0.042355932203389825|\n|036e5fae1ad94e66a...|0.042355932203389825|0.042355932203389825|\n+--------------------+--------------------+--------------------+\nonly showing top 5 rows\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"10. 在sample上新添加一列ma15_price，计算sample中每个销售渠道的15天移动平均消费总额。展示sample的前5行。预期输出如下，数字不会完全一样。\n\n```\n+----------+--------------------+----------+--------------------+----------------+-------------------+\n|     t_dat|         customer_id|article_id|               price|sales_channel_id|         ma15_price|\n+----------+--------------------+----------+--------------------+----------------+-------------------+\n|2018-09-20|dc7e7d29a490da345...| 662980003|0.033881355932203386|               1| 0.1947966101694915|\n|2018-09-23|91334d6e58d43831b...| 425978006|0.022016949152542376|               1|0.23545762711864404|\n|2018-09-26|a969036a7d0a4058e...| 507909017| 0.02288135593220339|               1|0.26594915254237284|\n|2018-10-02|760a6b6da095ef2e6...| 630762001|0.008457627118644067|               1|0.29305084745762705|\n|2018-10-03|e44209348f5e5e017...| 624066007|0.022864406779661017|               1| 0.3099830508474576|\n+----------+--------------------+----------+--------------------+----------------+-------------------+\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"11. 为了看到每个顾客的消费能力，将customer重新连接到sample上。展示连接后的行数。预期输出如下，数字不会完全一样。\n\n```\n1615\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"12. 在完成11题时，有同学错误地把代码写成了sample.join(customer)。这样连接后的DataFrame有多少行？和正确的做法差异在哪里？","metadata":{}},{"cell_type":"code","source":"sample.join(customer).count()\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"13. 使用SQL语句查看sample的前五行。","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"14. 统计每个销售渠道的销售总额和交易总量，但不要改变sample的内容。预期输出如下，数字不会完全一样。\n```\n+----------------+-----------------+------------+\n|sales_channel_id|       sum(price)|count(price)|\n+----------------+-----------------+------------+\n|               1|10.81245762711864|         472|\n|               2|33.84950847457627|        1143|\n+----------------+-----------------+------------+\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"15.在sample的基础上新建一列year表示销售年度,一列month表示月份,一列day表示周几。展示前5行，预期输出如下，数字不会完全一样。\n```\n+--------------------+----------+----------+--------------------+----------------+-------------------+--------------------+--------------------+----+-----+---+\n|         customer_id|     t_dat|article_id|               price|sales_channel_id|         ma15_price|           sum_price|          mean_price|year|month|day|\n+--------------------+----------+----------+--------------------+----------------+-------------------+--------------------+--------------------+----+-----+---+\n|dc7e7d29a490da345...|2018-09-20| 662980003|0.033881355932203386|               1| 0.1947966101694915|0.033881355932203386|0.033881355932203386|2018|    9|  5|\n|91334d6e58d43831b...|2018-09-23| 425978006|0.022016949152542376|               1|0.23545762711864404|0.022016949152542376|0.022016949152542376|2018|    9|  1|\n|a969036a7d0a4058e...|2018-09-26| 507909017| 0.02288135593220339|               1|0.26594915254237284| 0.02288135593220339| 0.02288135593220339|2018|    9|  4|\n|760a6b6da095ef2e6...|2018-10-02| 630762001|0.008457627118644067|               1|0.29305084745762705|0.008457627118644067|0.008457627118644067|2018|   10|  3|\n|e44209348f5e5e017...|2018-10-03| 624066007|0.022864406779661017|               1| 0.3099830508474576|0.022864406779661017|0.022864406779661017|2018|   10|  4|\n+--------------------+----------+----------+--------------------+----------------+-------------------+--------------------+--------------------+----+-----+---+\nonly showing top 5 rows\n```\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"16. 按照年度、销售渠道两个维度统计销售总额。有几种方法可以实现？彼此有何不同？","metadata":{}},{"cell_type":"code","source":"#编程区域\n","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"17. 在sample上选出符合以下条件的行，统计行数。\n- （1）销售渠道为 1；\n- （2）销售年度为 2020；\n- （3）商品单价大于 0.03。\n\n用where, filter, spark.sql, 分别实现","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"18. 在sample数据表的基础上，新建一列布尔变量 flag，用于标识是否满足以下条件的交易记录：\n- （1）销售渠道为 1；\n- （2）销售年度为 2020；\n- （3）商品单价大于 0.03。\n\n随后，将flag更改为0-1变量，统计 flag = 1 的记录行数。","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"19. 查看18题的逻辑计划和物理计划。","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"20. 基于19题的物理计划回答以下问题\n    -  （1）18题的计算，是从哪个步骤开始的？\n    -  （2）上述计划涉及了练习中的哪些题目？为什么这些题目的代码也出现在18题的逻辑计划中了？\n    -  （3）上述计划中是都发生了shuffle？是否存在宽依赖转换？是否存在在依赖转换？是否存在行动操作？\n    -  （4）为什么进行了两次csv scan?","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"**思考题**\n\n在完成前序多道数据处理与特征构造练习后，第18题中仅对已有数据表执行了行数统计操作。然而，从 Spark 生成的 Physical Plan 可以观察到，执行计划仍然包含从原始 CSV 数据读取、过滤、聚合、窗口计算以及连接等多个上游步骤。\n\n请结合 Spark 的执行模型，解释为什么在已经完成多步处理的情况下，Spark 仍会从数据源开始构建并执行完整的计算计划。进一步说明这种设计是否必然意味着计算资源的浪费，以及 Spark 通过哪些机制来避免不必要的重复计算。","metadata":{}},{"cell_type":"markdown","source":"**参考答案**\n\n在 Spark 中，DataFrame 的各类转换操作（如 filter、withColumn、groupBy、join 等）采用惰性求值机制。在代码执行过程中，这些操作并不会立即触发实际计算，而只是不断构建数据的**血统关系（lineage）**与逻辑执行计划。只有当遇到行动操作（action，如 count()、show() 等）时，Spark 才会基于当前完整的 lineage 生成物理执行计划并真正开始计算。\n\n因此，即使在第 18 题中仅执行了简单的行数统计操作，Spark 仍需要从最初的数据源出发，将所有与该结果相关的上游转换操作统一纳入执行计划。这一过程并非重复执行之前的计算步骤，而是确保 Spark 能够在全局范围内对执行过程进行优化，包括算子重排、过滤条件下推以及连接策略选择等。\n\n从物理执行计划可以观察到，部分中间结果被缓存为 InMemoryRelation，后续计算通过 InMemoryTableScan 直接从内存中读取数据，而非重复从 CSV 文件中加载。这表明 Spark 并未盲目重复 IO 操作，而是在保证计算正确性的前提下，尽可能复用已缓存的数据。此外，Catalyst 优化器与自适应查询执行（Adaptive Query Execution, AQE）机制也会根据运行时统计信息动态调整执行策略，从而减少不必要的计算和数据移动。\n\n综上所述，Spark 在第 18 题中看似“从头开始”的执行方式，本质上是惰性求值与全局优化策略共同作用的结果。这种设计并不必然导致计算资源浪费，反而通过完整的执行计划与缓存机制，在保证结果一致性的同时提升了整体执行效率。","metadata":{}},{"cell_type":"markdown","source":"# 第二部分：特征工程\n第二部分我们进入特征工程（Feature Engineering）：目标是把“原始业务字段”加工成机器学习模型能够直接学习的label（标签）和features（特征向量），从而构建一个可复用的预测 Pipeline，用来判断某次交易更可能发生在何种渠道（例如线上/线下）。之所以要做特征工程，是因为模型并不理解“customer_id、商品名称、类别、价格、年龄”等原始字段背后的业务含义；我们必须把这些信息转换成结构化、可计算、可比较的数值表示，让模型能够从历史数据中学习规律并对新样本进行预测。本部分将完成三件事：先通过抽样与多表连接把用户表、商品表与交易表整合成统一训练数据；再把类别型、数值型、文本型字段分别做编码、填补与离散化等处理；最后把这些处理步骤组织成一个稳定的特征工程 Pipeline，得到统一的 label 与 features，为后续的模型训练与调优打下基础。","metadata":{}},{"cell_type":"markdown","source":"21. 以学号作为随机种子，从 transactions_train.csv 中随机抽取 0.0001 比例的数据，命名为 data。\n在此基础上：\n\n- 使用 LEFT JOIN 的方式，将 articles.csv 连接到 data 上（按 article_id）；\n\n- 使用 INNER JOIN 的方式，将 customers.csv 连接到 data 上（按 customer_id）；\n\n- 将 data 聚合成8个分区，并缓存在内存中\n\n查看最终得到的 data 包含哪些字段（列）。\n\n预期输出如下\n\n```\nroot\n |-- customer_id: string (nullable = true)\n |-- article_id: integer (nullable = true)\n |-- t_dat: date (nullable = true)\n |-- price: double (nullable = true)\n |-- sales_channel_id: integer (nullable = true)\n |-- product_code: integer (nullable = true)\n |-- prod_name: string (nullable = true)\n |-- product_type_no: integer (nullable = true)\n |-- product_type_name: string (nullable = true)\n |-- product_group_name: string (nullable = true)\n |-- graphical_appearance_no: integer (nullable = true)\n |-- graphical_appearance_name: string (nullable = true)\n |-- colour_group_code: integer (nullable = true)\n |-- colour_group_name: string (nullable = true)\n |-- perceived_colour_value_id: integer (nullable = true)\n |-- perceived_colour_value_name: string (nullable = true)\n |-- perceived_colour_master_id: integer (nullable = true)\n |-- perceived_colour_master_name: string (nullable = true)\n |-- department_no: integer (nullable = true)\n |-- department_name: string (nullable = true)\n |-- index_code: string (nullable = true)\n |-- index_name: string (nullable = true)\n |-- index_group_no: integer (nullable = true)\n |-- index_group_name: string (nullable = true)\n |-- section_no: integer (nullable = true)\n |-- section_name: string (nullable = true)\n |-- garment_group_no: integer (nullable = true)\n |-- garment_group_name: string (nullable = true)\n |-- detail_desc: string (nullable = true)\n |-- FN: double (nullable = true)\n |-- Active: double (nullable = true)\n |-- club_member_status: string (nullable = true)\n |-- fashion_news_frequency: string (nullable = true)\n |-- age: integer (nullable = true)\n |-- postal_code: string (nullable = true)\n```\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"\n22. **思考题**\n\n为了便于后续多次特征工程与建模操作，我们在完成数据抽样与多表连接后，对中间结果 `data` 显式调用了 `cache()` 进行缓存处理。与此同时，Spark 本身采用惰性求值与全局优化机制，在行动操作触发时会基于完整 lineage 统一生成并优化执行计划。\n\n请结合 Spark 的执行模型，思考并回答以下问题：\n\n1. 在 DataFrame 上显式使用 `cache()`，是否会与 Spark 的全局优化理念产生冲突？\n2. `cache()` 在 Spark 中的实际作用是什么，它会在何时、以何种方式生效？\n3. 在什么样的使用场景下，对中间结果进行缓存是合理且有益的？在什么情况下，缓存反而可能造成资源浪费？\n","metadata":{"execution":{"iopub.status.busy":"2025-12-18T05:18:32.486269Z","iopub.execute_input":"2025-12-18T05:18:32.486586Z","iopub.status.idle":"2025-12-18T05:18:32.493688Z","shell.execute_reply.started":"2025-12-18T05:18:32.486561Z","shell.execute_reply":"2025-12-18T05:18:32.492475Z"}}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"23. 类别型特征编码（StringIndexer）：在第 21 题得到的 data 数据表基础上，选择商品表中的类别字段 product_type_name，使用 StringIndexer 将其转换为数值型特征 product_type_idx，将这个转换器命名为indexer便于后续使用。查看转换前该字段的取值示例。\n```\n+-----------------+----------------+\n|product_type_name|product_type_idx|\n+-----------------+----------------+\n|T-shirt          |3.0             |\n|Bracelet         |68.0            |\n|Trousers         |0.0             |\n|Necklace         |30.0            |\n|Dress            |1.0             |\n|Blouse           |4.0             |\n|Bikini top       |9.0             |\n|Hoodie           |15.0            |\n|Trousers         |0.0             |\n|Bra              |7.0             |\n+-----------------+----------------+\nonly showing top 10 rows\n\n```\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"24. 在第 23 题中，我们已使用 StringIndexer 将类别型字段product_type_name 转换为数值编码 product_type_idx。在此基础上，使用 OneHotEncoder 对 product_type_idx 进行独热编码，生成向量形式的特征 product_type_ohe。查看并对比编码前后的字段变化。\n```\n+-----------------+----------------+----------------+\n|product_type_name|product_type_idx|product_type_ohe|\n+-----------------+----------------+----------------+\n|T-shirt          |3.0             |(78,[3],[1.0])  |\n|Bracelet         |68.0            |(78,[68],[1.0]) |\n|Trousers         |0.0             |(78,[0],[1.0])  |\n|Necklace         |30.0            |(78,[30],[1.0]) |\n|Dress            |1.0             |(78,[1],[1.0])  |\n|Blouse           |4.0             |(78,[4],[1.0])  |\n|Bikini top       |9.0             |(78,[9],[1.0])  |\n|Hoodie           |15.0            |(78,[15],[1.0]) |\n|Trousers         |0.0             |(78,[0],[1.0])  |\n|Bra              |7.0             |(78,[7],[1.0])  |\n+-----------------+----------------+----------------+\nonly showing top 10 rows\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"25. 统计product_type_name的独特值数量。请思考为什么product_type_ohe只有76个元素，与这个数字不一样呢？","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"26. 查看data中的客户年龄字段是否存在缺失值。","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"27. 使用imputer转换器用中位数填充客户年龄的缺失值。检查填补前后的差异，预期输出如下：\n```\n+----+-----------+\n| age|age_imputed|\n+----+-----------+\n|NULL|         32|\n|NULL|         32|\n|NULL|         32|\n|NULL|         32|\n|NULL|         32|\n|NULL|         32|\n|NULL|         32|\n|NULL|         32|\n|NULL|         32|\n|NULL|         32|\n+----+-----------+\nonly showing top 10 rows\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"28. 数值型特征分位数离散化（QuantileDiscretizer）：在 data 数据表中，商品价格 price 为连续型数值特征，且分布高度不均衡。本题要求使用 QuantileDiscretizer 按分位数对 price 进行分组，将其离散化为10个价格区间，并生成新的特征列 price_bin。查看转换前后的数据。示例如下：\n```\n+--------------------+---------+\n|               price|price_bin|\n+--------------------+---------+\n|0.013542372881355931|      2.0|\n|0.012966101694915255|      1.0|\n|0.020322033898305086|      4.0|\n|0.012186440677966101|      1.0|\n|0.021338983050847457|      4.0|\n|0.022016949152542376|      4.0|\n| 0.03049152542372881|      6.0|\n|0.027101694915254236|      6.0|\n| 0.04744067796610169|      8.0|\n| 0.03049152542372881|      6.0|\n+--------------------+---------+\nonly showing top 10 rows\n```\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"29. 分词（Tokenizer）：在 data 数据表中，商品名称字段 prod_name 为文本型特征，无法直接用于模型训练。本题要求使用 Tokenizer 对该字段进行分词处理，将商品名称拆分为若干词项，并生成新的特征列 prod_tokens。对比分词前后字段的变化。示例如下：\n```\n+-----------------------------+-----------------------------------+\n|prod_name                    |prod_tokens                        |\n+-----------------------------+-----------------------------------+\n|TAZ LONG RN T-SHIRT          |[taz, long, rn, t-shirt]           |\n|Bracelet Lincoln Italy       |[bracelet, lincoln, italy]         |\n|Summer pants (Saigon)        |[summer, pants, (saigon)]          |\n|Flirty Trams necklace        |[flirty, trams, necklace]          |\n|Singo.                       |[singo.]                           |\n|Darling                      |[darling]                          |\n|Knot Bitter Top              |[knot, bitter, top]                |\n|Disa jacket                  |[disa, jacket]                     |\n|Guela jaquard                |[guela, jaquard]                   |\n|Margarita Push w. Back Detail|[margarita, push, w., back, detail]|\n+-----------------------------+-----------------------------------+\n```","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"30. 在上一题中，我们已使用 Tokenizer 对商品名称字段 prod_name 进行了分词处理，得到了词项列表 prod_tokens。本题要求进一步使用 StopWordsRemover 对分词结果进行停用词过滤，并将 Tokenizer 与 StopWordsRemover 组合为一个 Pipeline，生成新的特征列 prod_tokens_clean。","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"# 第三部分：机器学习\n\n## 31. **综合题：基于 Spark ML Pipeline 的特征工程与分类建模**\n\n在前序实验中，我们已经完成了交易数据 `data` 的构建（包括从交易表中抽样，并与商品表和客户表进行连接）。在此基础上，请使用 **Spark ML Pipeline** 构建一个完整的机器学习流程，对销售渠道进行二分类建模。\n\n本题要求你在 **不使用中间手工 DataFrame 转换** 的前提下，将以下步骤全部纳入同一个 Pipeline 中完成：\n\n1. **时间特征构造**\n   使用 `SQLTransformer`，从交易日期字段 `t_dat` 中提取：\n\n   * 交易月份（month）\n   * 星期几（day of week）\n\n     同时，将 `sales_channel_id` 转换为二分类标签列 `label`。\n\n2. **价格特征离散化**\n   使用 `QuantileDiscretizer`，将连续变量 `price` 按分位数划分为 **5 个价格区间**，生成离散变量 `price_bin`。\n\n3. **文本特征表示**\n   以商品名称字段 `prod_name` 为输入，完成以下文本处理流程：\n\n   * 分词（Tokenizer）\n   * 停用词过滤（StopWordsRemover）\n   * 使用 `Word2Vec` 将文本表示为定长向量特征\n\n4. **缺失值处理**\n   对客户年龄字段 `age` 的缺失值进行填补，生成新的数值型特征。\n\n5. **特征组装与缩放**\n   将上述得到的时间特征、价格特征、文本特征和年龄特征合并为统一的特征向量 `features`。\n   请思考是否有必要对特征进行标准化，并在 Pipeline 中给出你的处理方式。\n\n6. **模型构建与调参**\n   使用 **Logistic Regression** 作为分类模型，并通过 **参数网格（ParamGrid）** 与 **交叉验证（CrossValidator）** 对模型参数进行调优。\n\n\n打印你的最优模型的参数与AUC\n\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 31. 在本实验中，模型训练和参数搜索需要多次重复执行相同的数据处理与特征构建步骤。请解释 Spark 的计算特点为什么能够提高这类机器学习任务的运行效率。","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"## 32. 传统 MapReduce 主要用于一次性批处理任务，而 Spark 常被用于机器学习等迭代计算场景。请从计算过程和中间结果处理方式的角度，比较 Spark 与 MapReduce 的主要差异。","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"## 33. 在本实验中我们使用 Spark MLlib 完成模型训练与参数搜索。请从数据规模与计算方式的角度，比较 Spark MLlib 与 scikit-learn 的主要差异。","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"## 34. 在使用 CrossValidator 进行参数搜索时，Spark 会发生什么样的计算过程？为什么这种过程通常比较耗时？\n    ","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"## 35. 在 CrossValidator 中，由于需要多次重复执行模型训练，计算开销往往较大。请结合机器学习实践，说明可以从哪些方面降低 CrossValidator 的计算成本。","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"# 第四部分：推荐系统\n第四部分我们进入 **推荐系统（Recommender Systems）** 模块：在前面完成数据读取与 Spark 执行机制理解的基础上，我们将把 H&M 交易明细转换为适合推荐建模的“用户–物品–交互强度”结构，并使用  **ALS（交替最小二乘）** 训练一个协同过滤模型，为每位用户生成个性化商品推荐。本部分之所以要先新建 SparkSession，是为了让实验在一个干净、可控的运行环境中进行，避免前序题目残留的缓存、配置或临时视图影响结果；同时我们会显式设置 shuffle 分区数与广播阈值，便于观察任务并行度与性能表现。随后，你将经历一条完整的推荐建模流程：从原始交易数据抽样与缓存、到用户–物品层聚合与再抽样以降低稀疏性，再到通过 Pipeline + CrossValidator 进行参数调优，并在测试集上用 RMSE 评估模型质量。通过这一流程，你不仅要学会“把推荐跑出来”，更要理解推荐系统为何需要这种数据结构、ALS 学到的分数代表什么、以及 Spark 在分布式机器学习任务中如何调度与执行。","metadata":{}},{"cell_type":"code","source":"# 新建一个SparkSession,避免前序题目的干扰\nspark.stop()","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"from pyspark.sql import SparkSession\nspark = SparkSession.builder.config(\"spark.sql.shuffle.partitions\", 4\n                                   ).config(\"spark.sql.autoBroadcastJoinThreshold\",  -1\n                                           ).getOrCreate()","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# 查看计算机和Spark性能\nimport os\nprint(os.cpu_count())\nspark.sparkContext.defaultParallelism\nspark.sparkContext.getConf().getAll()","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 读取原始数据","metadata":{}},{"cell_type":"code","source":"%%time\nfrom pyspark.ml import Pipeline\nfrom pyspark.ml.feature import SQLTransformer, StringIndexer\nfrom pyspark.ml.recommendation import ALS\nfrom pyspark.ml.tuning import CrossValidator, ParamGridBuilder\nfrom pyspark.ml.evaluation import RegressionEvaluator\nfrom pyspark.sql import functions as F\n\n\n# 原始交易数据（减少分区 + cache + 抽样（10%））\ndata = spark.read.csv(\n    \"/kaggle/input/h-and-m-personalized-fashion-recommendations/transactions_train.csv\",\n    header=True,\n    inferSchema=True\n).coalesce(4).sample(0.1,1234).cache()\ndata.count()","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"code","source":"%%time\n# 确认分区与缓存情况\nprint(data.rdd.getNumPartitions())\nprint(data.count())","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"\n## 实验解析：Spark 中数据读取、分区与执行过程的理解\n\n在本实验中，我们使用如下代码对 H&M 交易数据进行抽样与缓存，并通过 `count()` 触发计算：\n\n```python\ndata = spark.read.csv(\n    \"transactions_train.csv\",\n    header=True,\n    inferSchema=True\n).coalesce(4).sample(0.1, 1234).cache()\n\ndata.count()\n```\n\n尽管代码看似简洁，但其背后涉及 Spark 的 **懒执行机制、分区策略、Stage 划分以及并行调度模型**。下面结合运行过程进行说明。\n\n---\n\n### 一、懒执行（Lazy Evaluation）与 Action 触发\n\n在 Spark 中，大多数 DataFrame 操作（如 `read`、`coalesce`、`sample`、`cache`）都是**转换（transformation）**，并不会立即执行计算。\n这些操作只是构建了一条**逻辑执行计划（lineage）**。\n\n真正的计算发生在遇到 **行动操作（action）** 时，例如：\n\n* `count()`\n* `show()`\n* `collect()`\n\n因此，本例中 **真正触发数据读取、抽样和缓存的操作是 `data.count()`**。\n\n---\n\n### 二、为什么会出现两个 Stage？\n\n运行过程中可以观察到类似如下输出：\n\n```\n[Stage 1: ... / 26]\n[Stage 2: ... / 4]\n```\n\n这反映了 Spark 在物理执行阶段对任务的拆分方式。\n\n#### 1️⃣ Stage 1：CSV 文件读取阶段（26 个任务）\n\n* Spark 在读取 CSV 文件时，会根据文件大小和内部参数（如文件切分大小）将输入数据划分为多个**输入分区（input partitions）**\n* 在本实验环境中，该 CSV 文件被切分为 **26 个分区**\n* 因此，Stage 27 包含 **26 个 task**，负责从磁盘读取数据\n\n> 注意：**这个分区数与 CPU 核数无关**，而是由 Spark 的文件读取策略决定。\n\n#### 2️⃣ Stage 2：`coalesce(4)` 之后的下游阶段（4 个任务）\n\n* `coalesce(4)` 是一个**窄依赖（narrow dependency）操作**，不会触发 shuffle\n* 但在物理执行中，Spark 仍然需要将上游的 26 个分区合并为 4 个分区\n* 因此，形成了第二个 Stage（Stage 28），其 task 数为 4\n\n该 Stage 同时完成以下工作：\n\n* 抽样（`sample(0.1)`）\n* 缓存数据（`cache()`）\n* 统计行数（`count()`）\n\n---\n\n### 三、任务并行度与 `(x + y) / n` 的含义\n\n运行日志中常见如下格式：\n\n```\n(8 + 4) / 26\n```\n\n其含义是：\n\n* `8`：已经完成的 task 数\n* `4`：当前正在运行的 task 数\n* `26`：该 Stage 的 task 总数\n\n在本实验中，Spark 运行于 `local[*]` 模式，且默认并行度为 4，因此：\n\n> **任意时刻最多只能并行执行 4 个 task，其余任务需要排队等待**\n\n---\n\n### 四、如何理解 CPU 使用率接近 400%？\n\n在 Kaggle Notebook 的资源监控中，可以观察到：\n\n* CPU 使用率 ≈ 393%\n\n这并不是异常，而是正常现象：\n\n* CPU 使用率是**按单核 100% 计量**\n* 393% ≈ 3.93 个 CPU 核心被同时使用\n* 与 Spark 在 `local[*]` 模式下最多使用 4 核完全一致\n\n---\n\n### 五、为什么第一次 `count()` 耗时较长？\n\n这是因为：\n\n* 第一次 `count()` 需要：\n\n  1. 从磁盘读取 CSV\n  2. 执行抽样\n  3. 将结果写入缓存\n  4. 同时统计行数\n\n如果随后再次执行：\n\n```python\ndata.count()\n```\n\n通常会明显更快，因为数据已经缓存在内存中，无需再次从磁盘读取。\n\n---\n\n### 六、小结\n\n通过这一简单示例，我们可以总结出 Spark 执行过程中的几个关键特性：\n\n* Spark 采用**懒执行模型**，只有在 action 出现时才真正计算\n* 文件读取阶段的分区数由 **文件切分策略** 决定，而非 CPU 核数\n* `coalesce` 会影响下游分区数，从而影响 task 数量\n* 在 `local[*]` 模式下，任务会按 CPU 核数排队执行\n* 第一次 action 往往承担“数据加载 + 计算 + 缓存”的全部开销\n\n理解这些机制，有助于在后续实验中 **正确解读 Stage、Task、并行度和性能表现**，而不是仅凭代码表面判断程序行为。","metadata":{}},{"cell_type":"markdown","source":"## 用户-物品-购买数量数据重构","metadata":{}},{"cell_type":"code","source":"%%time\n# 生成购买数量行式数据用于ALS，在用户-物品层再次抽样，减少对稀疏性的影响。\nsqlT = SQLTransformer().setStatement(\"\"\"\nWITH ui AS (\n  SELECT customer_id,\n         article_id,\n         COUNT(*) AS count\n  FROM __THIS__\n  GROUP BY customer_id, article_id\n)\nSELECT *\nFROM ui TABLESAMPLE (1 PERCENT)\n\"\"\")\n\nui = sqlT.transform(data).cache()\nprint(ui.count())\nui.show(5, truncate=False)","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 实验解析：基于 SQLTransformer 的用户–物品聚合、抽样与执行机制\n\n在推荐系统实验中，ALS 模型通常需要以 **用户–物品–交互强度** 作为输入。本实验通过 Spark SQL 对交易明细进行聚合，并在用户–物品层面进行抽样，以构建适合 ALS 训练的数据格式。同时，本实验也用于观察 **Spark 在不同执行配置下的任务划分与执行行为**。\n\n---\n\n### 一、用户–物品层数据的构造逻辑\n\n实验中使用如下 `SQLTransformer`：\n\n```python\nsqlT = SQLTransformer().setStatement(\"\"\"\nWITH ui AS (\n  SELECT customer_id,\n         article_id,\n         COUNT(*) AS count\n  FROM __THIS__\n  GROUP BY customer_id, article_id\n)\nSELECT *\nFROM ui TABLESAMPLE (1 PERCENT)\n\"\"\")\n```\n\n该 SQL 包含两个关键步骤：\n\n### 1️⃣ 用户–物品聚合（GROUP BY）\n\n将原始交易明细（每行代表一次购买）聚合为：\n\n* 每个用户（`customer_id`）\n* 每个商品（`article_id`）\n* 对应的购买次数（`count`）\n\n该结果正是 ALS 推荐模型所需的典型输入形式。\n\n### 2️⃣ 在用户–物品层进行抽样（TABLESAMPLE）\n\n抽样发生在聚合之后，而非交易明细层：\n\n* 抽取的是一部分 `(user, item)` 关系\n* 而不是一部分原始交易记录\n\n这种做法在教学和实验中有两个目的：\n\n* 降低数据规模，缩短实验运行时间\n* 在尽量保留用户–物品交互结构的前提下，快速构建可运行的 ALS 示例\n\n---\n\n### 二、cache 与 action：什么时候真正执行计算？\n\n```python\nui = sqlT.transform(data).cache()\nprint(ui.count())\nui.show(5, truncate=False)\n```\n\n在 Spark 中：\n\n* `transform` 和 `cache` 都属于**懒执行（lazy evaluation）**操作\n* 它们只是在构建执行计划，并不会立即触发计算\n\n真正的执行发生在以下 **action** 被调用时：\n\n* `count()`\n* `show()`\n\n此时 Spark 才会：\n\n* 从上游 DataFrame `data` 读取数据\n* 执行 `GROUP BY` 聚合\n* 对聚合结果进行抽样\n* 将结果写入缓存\n* 返回行数或展示样本数据\n\n---\n\n### 三、Stage 数量与 shuffle 分区的来源（重要）\n\n在实验运行过程中，可以观察到类似如下输出：\n\n```text\n[Stage … : … / N]\n```\n\n其中 `N` 表示该 Stage 中的任务（task）总数。\n该任务数量来源于 **Spark SQL 的 shuffle 分区机制**。\n\n### 1️⃣ GROUP BY 与 shuffle\n\n* `GROUP BY` 属于 **宽依赖（wide dependency）** 操作\n* 会触发 shuffle，将数据按 key 重新分发\n* shuffle 之后的数据会被划分为若干分区，并由多个 task 并行处理\n\n### 2️⃣ 默认行为（未显式设置时）\n\n在未进行额外配置时，Spark SQL 使用默认参数：\n\n```text\nspark.sql.shuffle.partitions = 200\n```\n\n因此，`GROUP BY` 后的 shuffle 结果会被划分为 **200 个 task**，即日志中常见的：\n\n```text\n[Stage … : … / 200]\n```\n\n在 `local[*]` 模式下：\n\n* 计算机只有 4 个 CPU 核\n* 任意时刻最多并行执行 4 个 task\n* 其余 task 需要排队执行\n\n---\n\n### 3️⃣ 本实验中的配置调整\n\n在本实验中，我们在创建 `SparkSession` 时显式设置了：\n\n```python\n.config(\"spark.sql.shuffle.partitions\", 4)\n```\n\n因此，`GROUP BY` 等 shuffle 操作不再生成 200 个任务，而是生成 **4 个任务**：\n\n```text\n[Stage … : … / 4]\n```\n\n需要注意的是：\n\n> **shuffle 分区数的变化并不改变计算逻辑，只影响任务切分粒度与调度方式。**\n\n---\n\n### 四、关于 RowBasedKeyValueBatch 的 warning\n\n运行过程中可能出现如下警告：\n\n```text\nWARN RowBasedKeyValueBatch: Calling spill() on RowBasedKeyValueBatch. Will not spill but return 0.\n```\n\n该警告可以这样理解：\n\n* Spark 在执行聚合时，会在内存中维护中间的 key–value 聚合结构\n* 当 Spark 认为内存压力可能增大时，会尝试触发 *spill*（将部分中间结果写入磁盘）\n* 但当前使用的内部数据结构（`RowBasedKeyValueBatch`）不支持该方式的 spill\n* 因此 Spark 发出 warning 后继续执行\n\n需要强调的是：\n\n* 这是一个 **运行时提示（warning）**\n* 并不表示任务失败\n* 在内存充足、作业正常完成的情况下，通常不影响结果正确性\n\n---\n\n### 五、Wall time 与 CPU time 的差异\n\n实验中可以看到：\n\n* CPU time 非常小\n* Wall time 却达到数秒\n\n这是因为：\n\n* `%%time` 中的 CPU time 主要统计的是 **Python Notebook 进程**\n* Spark 的实际计算发生在 JVM 后台进程中\n* 因此真实计算开销体现在 Wall time 与系统资源监控中\n\n---\n\n### 六、小结\n\n通过本实验可以总结出以下关键点：\n\n* Spark SQL 中的 `GROUP BY` 会触发 shuffle，并根据 shuffle 分区数生成对应数量的 task\n* shuffle 分区数默认是 200，但可通过 `spark.sql.shuffle.partitions` 显式调整\n* `SQLTransformer` 与 `cache()` 本身不会触发计算，必须通过 action 才会执行\n* 在用户–物品层进行抽样，有助于在推荐系统实验中快速构建可运行样本\n* 聚合阶段的内存管理可能产生 warning，但通常不影响计算结果\n* Spark 的执行过程以 Stage 和 Task 为基本单位，其并行度与运行环境配置密切相关\n\n理解这些执行细节，有助于在后续推荐系统与大规模机器学习实验中，**正确解读 Spark 的运行日志与性能表现，而不仅停留在代码层面**。\n\n---","metadata":{}},{"cell_type":"markdown","source":"## ALS模型训练与调优\n## 36. 综合题：ALS 模型训练与调优（Pipeline + CrossValidator）\n\n### 已知条件\n\n你已完成“原始数据读取 + 用户–物品–购买次数重构”，得到了 DataFrame：`ui`，其至少包含三列：\n\n* `customer_id`（string）\n* `article_id`（string）\n* `count`（int/float，购买次数）\n\n---\n\n### 任务\n\n请基于 `ui` 完成 ALS 推荐模型的**训练、交叉验证调优、测试集评估与推荐输出**，并在关键步骤写明必要注释。\n\n---\n\n### 要求\n\n#### A. 数据划分\n\n1. 按照0.7: 0.3, seed=42将 `ui` 划分为 `train` 和 `test`。\n\n#### B. 特征处理与建模管道（Pipeline）\n\n2. 分别对 `customer_id` 与 `article_id` 做 `StringIndexer`：\n\n   * `customer_id` → `customer_id_index`\n   * `article_id` → `article_id_index`\n     且设置 `handleInvalid=\"skip\"`。\n3. 构建 ALS 模型，要求参数如下：\n\n   * `userCol=`\n   * `itemCol=`\n   * `ratingCol=`\n   * `nonnegative=True`\n   * `coldStartStrategy=\"drop\"`\n   * `seed=42`\n4. 将两个 indexer 和 ALS 组成 `Pipeline(stages=[...])`。\n\n#### C. 参数网格与交叉验证\n\n5. 构建参数网格 `ParamGridBuilder`，搜索空间为：\n\n   * `rank ∈ {20, 25}`\n   * `maxIter ∈ {5, 10}`\n   * `regParam ∈ {0.05, 0.10}`\n6. 构建 `RegressionEvaluator`：\n\n   * `metricName=`\n   * `labelCol=`\n   * `predictionCol=\"prediction\"`\n7. 构建 `CrossValidator`：\n\n   * `estimator=pipe`\n   * `estimatorParamMaps=param_grid`\n   * `evaluator=evaluator`\n   * `numFolds=3`\n   * `parallelism=2`\n\n#### D. 最优模型选择与测试集评估\n\n8. 在 `train` 上训练得到 `cvModel`，并取出：\n\n   * `best_pipe = cvModel.bestModel`\n   * `best_als = best_pipe.stages[-1]`\n9. 用 `best_pipe.transform(test)` 得到预测结果 `pred`，并计算测试集 RMSE。\n10. 打印输出（必须包含以下字段）：\n\n* `RMSE = ...`\n* `Rank`、`MaxIter`、`RegParam`（三者缺一不得分）\n\n#### E. 推荐输出\n\n11. 使用最优 ALS 模型对所有用户生成 Top-3 推荐：\n\n* `best_als.recommendForAllUsers(3)`\n\n12. 展示前 5 行，要求：`show(5, truncate=False)`。\n\n---\n\n### 代码骨架（请补全所有 TODO）\n\n```python\n%%time\nfrom pyspark.ml import Pipeline\nfrom pyspark.ml.feature import StringIndexer\nfrom pyspark.ml.recommendation import ALS\nfrom pyspark.ml.tuning import ParamGridBuilder, CrossValidator\nfrom pyspark.ml.evaluation import RegressionEvaluator\n\n# =========================================================\n# A. train/test 划分\n# =========================================================\n# TODO: \ntrain, test = ui.randomSplit(______________, seed=______)\n\n# =========================================================\n# B. CV 管道：Indexer + ALS\n# =========================================================\n# TODO:\nuser_indexer = StringIndexer(\n    inputCol=\"__________\", outputCol=\"__________________\", handleInvalid=\"____\"\n)\n\n# TODO: \nitem_indexer = StringIndexer(\n    inputCol=\"__________\", outputCol=\"__________________\", handleInvalid=\"____\"\n)\n\n# TODO: 构建 ALS（按题目要求设置参数）\nals = ALS(\n    userCol=\"__________________\",\n    itemCol=\"__________________\",\n    ratingCol=\"_____\",\n    nonnegative=_____,\n    coldStartStrategy=\"_____\",\n    seed=_____\n)\n\n# TODO: Pipeline 组装\npipe = Pipeline(stages=[______________________________])\n\n# =========================================================\n# C. ParamGrid + CrossValidator\n# =========================================================\n# TODO: \nparam_grid = (ParamGridBuilder()\n    .addGrid(als.______, [____, ____])\n    .addGrid(als.______, [____, ____])\n    .addGrid(als.______, [____, ____])\n    .build()\n)\n\n# TODO:\nevaluator = RegressionEvaluator(\n    metricName=\"_____\", labelCol=\"_____\", predictionCol=\"_____\"\n)\n\n# TODO: \ncv = CrossValidator(\n    estimator=_____,\n    estimatorParamMaps=_____,\n    evaluator=_____,\n    numFolds=_____,\n    parallelism=_____\n)\n\n# =========================================================\n# D. 训练 + 评估 + 输出最优参数\n# =========================================================\ncvModel = cv.fit(_____)\nbest_pipe = ______________________\nbest_als = ________________________\n\npred = best_pipe.transform(_____)\nrmse = evaluator.evaluate(pred)\n\nprint(\"RMSE =\", rmse)\nprint(\"=== Best ALS Params ===\")\nprint(\"Rank     :\", ____________________)\nprint(\"MaxIter  :\", ____________________)\nprint(\"RegParam :\", ____________________)\n\n# =========================================================\n# E. 推荐示例\n# =========================================================\nbest_als.recommendForAllUsers(____).show(____, truncate=____)\n```\n\n---\n\n### 提醒：常见失误\n\n* 忘记用 `Pipeline`，直接把 ALS 丢给 CV（会导致字符串列无法进入 ALS）\n* `handleInvalid` 未设为 `\"skip\"`，导致某些 fold 报错\n* 评估时没用 `best_pipe.transform(test)`（只用 ALS transform 会缺 indexer）\n* 打印最优参数不全或打印错对象（不是 `best_als`）\n\n---\n","metadata":{}},{"cell_type":"code","source":"","metadata":{"trusted":true},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"## 37. 上题运行过程中，Spark 多次输出如下警告：`WARN DAGScheduler: Broadcasting large task binary with size ~2 MB`。请结合 AI 助手的解释，说明该警告产生的原因，并阐明 Spark 在任务调度时“广播 task binary”这一设计的含义。该警告是否会影响模型训练结果？为什么？","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"**参考答案**\n\n在本实验运行过程中，Spark 多次输出\n`WARN DAGScheduler: Broadcasting large task binary with size ~2 MB`\n这一警告，说明 Spark 在向 Executor 分发任务时，需要广播一个体积较大的 **task binary（任务执行体）**。\n\n该警告产生的原因在于：Spark 在执行任务时，并不是逐条向 Executor 发送计算指令，而是将**完整的任务执行逻辑**进行序列化后整体广播给各个 Executor。在本实验中，任务中包含了较为复杂的机器学习 Pipeline，例如 SQLTransformer、StringIndexer、ALS 模型以及 CrossValidator 的参数配置等，这些组件及其内部状态都会被一并打包进 task binary，从而导致其体积达到数 MB。\n\nSpark 采用“广播 task binary”的设计，是为了减少 Driver 与 Executor 之间的频繁通信，使 Executor 能够在本地独立完成任务执行，提高分布式计算的效率。这种设计特别适合包含迭代计算和复杂逻辑的机器学习任务。\n\n该警告本身不会影响模型训练结果的正确性。它仅是 Spark 对任务调度开销的性能提示，表明任务结构较为复杂、序列化体积较大。在当前实验规模和环境下，该警告不会导致计算失败或结果错误，模型仍能正常完成训练与评估。\n","metadata":{}},{"cell_type":"markdown","source":"## 38. 本实验使用的 ALS 模型在建模过程中并未显式使用商品的内容特征或用户属性特征，但仍然能够生成推荐结果。请结合代码说明：\n- 这种推荐系统的类型是什么？\n- ALS 是基于什么信息学习用户偏好和商品相似性的？\n- 这种信息为什么被称为“协同（Collaborative）”信息？","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"## 39.在本实验中，count（购买次数）被作为 ratingCol 输入 ALS 模型，并使用 RMSE 作为评估指标。请思考并回答：\n- 这种设定更接近于显式反馈推荐还是隐式反馈推荐？\n- 如果将 implicitPrefs=True 打开，但评估方式仍然使用 RMSE，会产生什么问题？","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}},{"cell_type":"markdown","source":"## 40. 假设你希望构建一个基于商品属性（如类别、价格区间、品牌）的推荐系统，或者一个同时结合用户行为与商品内容信息的推荐系统：\n- 你认为当前这套 ALS + Pipeline + CV 的框架需要在哪些方面进行调整？\n- 是否还能把它归类为“纯协同过滤”？为什么？","metadata":{}},{"cell_type":"markdown","source":"答题区域：\n\n-\n\n---","metadata":{}}]}