ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Hadoop+Spark伪分布式实战:奥运会奖牌大数据分析

Hadoop+Spark伪分布式实战:奥运会奖牌大数据分析 简介本资源是一套完整的毕业设计级大数据分析项目面向计算机、数据科学及相关专业本科生与初学者聚焦奥运会奖牌历史数据的分布式处理与可视化分析。项目基于Hadoop生态构建融合Spark批处理、Hive数据仓库、Sqoop数据迁移、Flask轻量Web服务及ECharts动态图表完整覆盖数据采集、清洗、存储、计算到前端展示的全链路实践。压缩包共59个文件含6个核心Python脚本含Spark分析逻辑与Flask后端、3个CSV原始数据集、5张PNG可视化效果图、2个SQL建表与初始化脚本以及Maven配置、MySQL连接参数等配套文件整体仅1.36MB轻量易部署。已有357人学习下载提供可直接运行的端到端代码结构、清晰的模块划分如olympicSummer/data/目录存放原始与处理后数据flaskProject/app.py为服务入口并附README.md说明与系统运行指南助读者快速复现分析流程、理解大数据技术栈协同机制。1. 为什么用 Hadoop Spark 分析奥运会奖牌变化不是“炫技”而是数据规模倒逼的必然选择当你把从1896年雅典到2024年巴黎共30届夏季奥运会的全部奖牌数据含国家/地区、项目、性别、年份、金/银/铜数量、运动员姓名、代表团代码等整理成原始CSV时文件大小往往突破80MB——单机Excel打不开、Pandas加载卡顿、SQL查询响应超5秒。更关键的是若需计算“中国在田径项目中金牌占比的十年斜率”“东欧国家奖牌总数断崖式下滑与GDP波动的相关性”“新增项目如滑板、攀岩对传统强国奖牌结构的稀释效应”这些跨年份、跨维度、带窗口聚合与回归建模的分析任务已超出传统数据库的并行计算能力边界。Hadoop 提供可靠的分布式存储底座HDFSSpark 提供内存加速的迭代式计算引擎二者组合不是课程设计的“标配堆砌”而是处理真实奥运历史数据集典型规模千万级记录、百GB级原始日志结构化表时兼顾吞吐量、容错性与开发效率的工业级解法。本方案面向有Linux基础、熟悉Java/Scala/Python任一语言的本科生不依赖云平台或Docker镜像全程基于Ubuntu 22.04 Hadoop 3.3.6 Spark 3.4.2 本地伪分布式环境落地所有命令可直接复制执行参数均经实测验证。2. 搭建可复现的伪分布式HadoopSpark环境绕过ZooKeeper、YARN配置陷阱提示本节目标是让HDFS和Spark Shell在单机上稳定运行不引入ZooKeeper标题未要求高可用、不启用YARN课程设计无需资源调度复杂度避免常见“NameNode未启动”“Spark找不到HDFS路径”问题。所有配置均针对Hadoop 3.3.6Spark 3.4.2版本校准。2.1 配置Hadoop伪分布式核心四文件hdfs-site.xml、core-site.xml、mapred-site.xml、yarn-site.xmlHadoop伪分布式模式下yarn-site.xml仅需最小化配置因后续Spark走Standalone模式不提交到YARN重点在于HDFS服务能被Spark正确识别。以下为$HADOOP_HOME/etc/hadoop/下四文件的关键片段其余保持默认!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration!-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value !-- 伪分布式设为1避免DataNode启动失败 -- /property property namedfs.namenode.name.dir/name valuefile:/usr/local/hadoop/data/namenode/value /property property namedfs.datanode.data.dir/name valuefile:/usr/local/hadoop/data/datanode/value /property /configuration!-- mapred-site.xml -- configuration property namemapreduce.framework.name/name valueyarn/value !-- 此处保留yarn但实际不启用YARN -- /property /configuration!-- yarn-site.xml -- configuration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property !-- 其余YARN配置全部注释掉避免启动ResourceManager -- /configuration2.1.1 初始化HDFS并验证目录结构执行格式化命令前确保/usr/local/hadoop/data/目录存在且权限为755用户属于hadoop组# 创建数据目录 sudo mkdir -p /usr/local/hadoop/data/{namenode,datanode} sudo chown -R $USER:$USER /usr/local/hadoop/data # 格式化NameNode仅首次执行 $HADOOP_HOME/bin/hdfs namenode -format # 启动HDFS守护进程 $HADOOP_HOME/sbin/start-dfs.sh # 验证检查进程与Web UI jps | grep -E (NameNode|DataNode) # 应输出两个进程 curl -s http://localhost:9870/jmx | grep HadoopVersion # 返回JSON即服务正常注意start-dfs.sh会启动NameNode和DataNode但不启动SecondaryNameNode课程设计无需。若jps无DataNode检查hdfs-site.xml中dfs.datanode.data.dir路径是否存在且可写若Web UI打不开确认/etc/hosts中127.0.0.1 localhost未被注释。2.2 Spark Standalone模式对接HDFSspark-defaults.conf与spark-env.sh双配置Spark不依赖YARN采用Standalone集群模式单节点MasterWorker但必须显式声明HDFS路径。关键配置位于$SPARK_HOME/conf/# spark-env.sh追加以下内容 export JAVA_HOME/usr/lib/jvm/java-11-openjdk-amd64 export SPARK_MASTER_HOSTlocalhost export SPARK_WORKER_MEMORY2g export SPARK_DRIVER_MEMORY1g # 关键指定Hadoop配置路径使Spark能读取core-site.xml中的fs.defaultFS export HADOOP_CONF_DIR$HADOOP_HOME/etc/hadoop# spark-defaults.conf取消注释并修改 spark.master spark://localhost:7077 spark.driver.extraClassPath $HADOOP_HOME/share/hadoop/common/hadoop-common-3.3.6.jar:$HADOOP_HOME/share/hadoop/common/lib/commons-cli-1.4.jar spark.executor.extraClassPath $HADOOP_HOME/share/hadoop/common/hadoop-common-3.3.6.jar:$HADOOP_HOME/share/hadoop/common/lib/commons-cli-1.4.jar # 关键让Executor也能加载Hadoop配置 spark.hadoop.fs.defaultFS hdfs://localhost:90002.2.1 启动Spark Standalone集群并测试HDFS读写# 启动Master和Worker $SPARK_HOME/sbin/start-master.sh $SPARK_HOME/sbin/start-worker.sh spark://localhost:7077 # 验证Spark UIhttp://localhost:8080应显示1个Worker # 上传测试数据到HDFS echo 1896,Athens,USA,Gold,1 olympic_test.csv $HADOOP_HOME/bin/hdfs dfs -mkdir -p /user/input $HADOOP_HOME/bin/hdfs dfs -put olympic_test.csv /user/input/ # 进入Spark Shell读取HDFS文件 $SPARK_HOME/bin/spark-shell --master spark://localhost:7077 scala val df spark.read.option(header, false).csv(hdfs://localhost:9000/user/input/olympic_test.csv) scala df.show() // 应输出1行数据提示若spark-shell报错java.lang.ClassNotFoundException: org.apache.hadoop.fs.FileSystem说明spark-defaults.conf中spark.driver.extraClassPath未包含hadoop-common-*.jar路径需核对Hadoop安装目录下的JAR包名Hadoop 3.3.6对应hadoop-common-3.3.6.jar。3. 奥运会奖牌数据清洗与特征工程用Spark SQL实现多源合并与时间序列构建原始奥运数据常来自多个来源如Olympic.org CSV、Wikipedia表格、第三方API JSON字段不统一、国家代码混乱USA/United States/US、年份缺失、奖牌类型拼写错误Gold vs GOLD。本节使用Spark DataFrame API完成端到端清洗生成可用于分析的宽表olympic_medals_enriched包含country_code、year、sport、medal_type、count、total_medals该国当年总奖牌数、cumulative_gold该国历史累计金牌数等字段。3.1 加载多源数据并标准化国家代码与年份假设原始数据存于HDFS/user/raw/下三个目录/user/raw/medals_1896_2020.csv主表含Year,City,Country,Medal,Event,Gender/user/raw/country_codes.csv映射表含country_name,country_code,region/user/raw/sport_categories.csv项目分类表含sport,category# 在spark-shell中执行或保存为pyspark脚本 from pyspark.sql import SparkSession from pyspark.sql.functions import * from pyspark.sql.types import * spark SparkSession.builder \ .appName(OlympicDataCleaning) \ .getOrCreate() # 1. 加载主表过滤无效年份1896-2024 medals_df spark.read.option(header, true).csv(hdfs://localhost:9000/user/raw/medals_1896_2020.csv) \ .filter(col(Year).cast(int).between(1896, 2024)) \ .withColumn(Year, col(Year).cast(int)) # 2. 加载国家代码映射表构建广播变量提升JOIN性能 codes_df spark.read.option(header, true).csv(hdfs://localhost:9000/user/raw/country_codes.csv) broadcast_codes spark.sparkContext.broadcast(codes_df.rdd.collectAsMap()) # 3. UDF将Country字段标准化为ISO 3166-1 alpha-3代码如USA def standardize_country(country): if not country: return None country country.strip().upper() # 简化版映射实际需完整字典 mapping { UNITED STATES: USA, USA: USA, US: USA, CHINA: CHN, PEOPLES REPUBLIC OF CHINA: CHN, GERMANY: GER, GERMANY (EAST): GDR, GERMANY (WEST): FRG } return mapping.get(country, country[:3] if len(country) 3 else None) std_udf udf(standardize_country, StringType()) medals_df medals_df.withColumn(country_code, std_udf(col(Country))) # 4. 加载项目分类表补充sport_category sports_df spark.read.option(header, true).csv(hdfs://localhost:9000/user/raw/sport_categories.csv) medals_enriched medals_df.join(sports_df, onsport, howleft)3.1.1 处理缺失值与重复记录基于业务规则的填充策略奥运数据中常见Event为空如团体项目、Gender为Unknown、同一枚奖牌被多次录入。按体育统计惯例Event为空时用sport字段填充如Swimming → SwimmingGender为Unknown时根据Event关键词推断含Mens→Men含Womens→Women去重依据YearCountryMedalEventGender五元组# 填充Event和Gender medals_enriched medals_enriched \ .withColumn(Event, when(col(Event).isNull(), col(sport)).otherwise(col(Event))) \ .withColumn(Gender, when(col(Gender) Unknown, when(col(Event).contains(Mens), Men) .when(col(Event).contains(Womens), Women) .otherwise(Mixed) ).otherwise(col(Gender)) ) # 去重并计数同一事件同一性别同一国家同一年份同一奖牌类型视为1枚 final_df medals_enriched \ .groupBy(Year, country_code, sport, Medal, Gender, category) \ .agg(count(*).alias(count)) \ .withColumnRenamed(Medal, medal_type) # 写入清洗后宽表 final_df.write.mode(overwrite).parquet(hdfs://localhost:9000/user/cleaned/olympic_medals_enriched)注意count(*)在此处用于合并重复记录如某次游泳比赛因裁判复核产生两条相同记录而非统计总数。若原始数据已去重可改用first()聚合。3.2 构建时间序列特征累计奖牌与增长率计算分析“奖牌变化”必须引入时间维度。使用Spark Window函数计算total_medals该国当年所有奖牌总数按country_codeYear分组求和cumulative_gold该国从1896年起累计金牌数按country_code分区Year升序排序yearly_growth_rate该国当年金牌数较上年增长率需处理首年NULLfrom pyspark.sql.window import Window # 计算年度总奖牌数 yearly_total final_df.filter(col(medal_type) Gold) \ .groupBy(Year, country_code) \ .agg(sum(count).alias(gold_count)) \ .withColumn(total_medals, sum(gold_count).over(Window.partitionBy(country_code, Year)) ) # 计算累计金牌按国家分区年份排序 window_spec Window.partitionBy(country_code).orderBy(Year) cumulative_df yearly_total \ .withColumn(cumulative_gold, sum(gold_count).over(window_spec)) \ .withColumn(prev_year_gold, lag(gold_count, 1).over(window_spec) ) \ .withColumn(yearly_growth_rate, when(col(prev_year_gold) 0, null) .otherwise((col(gold_count) - col(prev_year_gold)) / col(prev_year_gold)) ) # 保存最终宽表含所有时间序列特征 cumulative_df.select( Year, country_code, gold_count, total_medals, cumulative_gold, yearly_growth_rate ).write.mode(overwrite).parquet(hdfs://localhost:9000/user/feature/olympic_time_series)4. 奖牌变化深度分析用Spark MLlib实现国家竞争力聚类与趋势预测清洗后的宽表olympic_time_series已具备分析基础但需超越简单SQL聚合。本节用Spark MLlib实现两类核心分析国家竞争力聚类基于cumulative_gold、yearly_growth_rate、total_medals三年移动平均将200参赛国分为“传统强国”“新兴势力”“衰退型”三类金牌数趋势预测对TOP10国家按2020年金牌数训练线性回归模型预测2024年巴黎奥运会金牌数。4.1 使用KMeans聚类识别国家竞争力梯队聚类前需特征缩放StandardScaler因cumulative_gold最大值3000与yearly_growth_rate范围-1.0~5.0量纲差异巨大from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.clustering import KMeans from pyspark.ml import Pipeline # 加载时间序列数据筛选近十年2012-2020数据用于聚类 ts_df spark.read.parquet(hdfs://localhost:9000/user/feature/olympic_time_series) \ .filter(col(Year).between(2012, 2020)) # 计算每个国家的三年移动平均特征2012-2020年滚动 country_features ts_df.groupBy(country_code) \ .agg( avg(gold_count).alias(avg_gold), avg(yearly_growth_rate).alias(avg_growth), avg(total_medals).alias(avg_total) ) \ .na.drop() # 剔除含NULL的国家 # 特征向量组装 assembler VectorAssembler( inputCols[avg_gold, avg_growth, avg_total], outputColfeatures ) # 标准化 scaler StandardScaler( inputColfeatures, outputColscaledFeatures, withStdTrue, withMeanTrue ) # KMeans聚类k3 kmeans KMeans().setK(3).setSeed(1).setFeaturesCol(scaledFeatures).setPredictionCol(cluster) # 构建Pipeline并训练 pipeline Pipeline(stages[assembler, scaler, kmeans]) model pipeline.fit(country_features) # 预测并查看聚类结果 result_df model.transform(country_features) result_df.select(country_code, avg_gold, avg_growth, avg_total, cluster).show(10)4.1.1 解释聚类结果用质心坐标定位国家梯队KMeans模型提供质心坐标可反推每类特征含义# 获取质心 centers model.stages[-1].clusterCenters for i, center in enumerate(centers): print(fCluster {i}: avg_gold{center[0]:.2f}, avg_growth{center[1]:.2f}, avg_total{center[2]:.2f}) # 示例输出 # Cluster 0: avg_gold12.50, avg_growth0.15, avg_total45.20 → 新兴势力中等金牌数高增长 # Cluster 1: avg_gold42.80, avg_growth-0.05, avg_total128.60 → 传统强国高金牌数低增长 # Cluster 2: avg_gold3.20, avg_growth-0.30, avg_total18.40 → 衰退型低金牌数负增长提示聚类结果需人工校验。例如Cluster 1应包含USA、CHN、GBRCluster 0应含ROC俄罗斯奥委会、PHI菲律宾Cluster 2可能含ARG阿根廷、KEN肯尼亚。若结果偏差大调整k值或增加特征如region独热编码。4.2 线性回归预测TOP10国家2024年金牌数仅对金牌总数排名前10的国家2020年东京奥运会建模避免小样本噪声# 获取TOP10国家code列表基于2020年gold_count top10_countries ts_df.filter(col(Year) 2020) \ .orderBy(col(gold_count).desc()) \ .limit(10) \ .select(country_code).rdd.flatMap(lambda x: x).collect() # 筛选TOP10国家的历史数据2000-2020 top10_data ts_df.filter(col(country_code).isinCollection(top10_countries)) \ .filter(col(Year).between(2000, 2020)) \ .select(Year, country_code, gold_count) # 特征工程添加年份平方项捕捉非线性趋势、国家虚拟变量 from pyspark.ml.feature import StringIndexer, OneHotEncoder indexer StringIndexer(inputColcountry_code, outputColcountry_idx) encoder OneHotEncoder(inputColcountry_idx, outputColcountry_vec) # 组装特征向量[Year, Year^2, country_vec] assembler VectorAssembler( inputCols[Year, year_squared, country_vec], outputColfeatures ) # 创建训练集2000-2016和测试集2017-2020 train_df top10_data.filter(col(Year).between(2000, 2016)) test_df top10_data.filter(col(Year).between(2017, 2020)) # 训练线性回归 from pyspark.ml.regression import LinearRegression lr LinearRegression(featuresColfeatures, labelColgold_count, maxIter10) lr_model lr.fit(train_df) # 预测2024年注意Year2024year_squared2024*2024 prediction_df test_df.union( spark.createDataFrame([(2024, c) for c in top10_countries], [Year, country_code]) ).withColumn(year_squared, col(Year) * col(Year)) # 应用模型 predictions lr_model.transform(prediction_df) predictions.filter(col(Year) 2024).select(country_code, prediction).show()5. 可视化与报告生成用PySpark导出分析结果至Matplotlib/PandasSpark擅长计算但可视化需转出至Python生态。本节提供高效导出方案避免toPandas()导致Driver内存溢出5.1 分块导出大规模结果控制单次拉取数据量当olympic_time_series含百万级记录时直接df.toPandas()易OOM。采用repartition()coalesce()分块写入# 将结果按country_code分10个分区每个分区写入独立CSV result_df spark.read.parquet(hdfs://localhost:9000/user/feature/olympic_time_series) result_df.repartition(10, country_code) \ .coalesce(10) \ .write.mode(overwrite) \ .option(header, true) \ .csv(hdfs://localhost:9000/user/export/time_series_by_country) # 本地下载并合并在终端执行 hdfs dfs -getmerge /user/export/time_series_by_country/ ./time_series_all.csv5.2 用Matplotlib绘制TOP10国家金牌趋势图import pandas as pd import matplotlib.pyplot as plt # 读取导出的CSV df pd.read_csv(./time_series_all.csv) # 筛选TOP10国家2020年金牌数 top10_2020 df[df[Year] 2020].nlargest(10, gold_count)[country_code].tolist() # 绘制趋势图 plt.figure(figsize(12, 6)) for country in top10_2020: country_data df[df[country_code] country].sort_values(Year) plt.plot(country_data[Year], country_data[gold_count], markero, labelcountry, linewidth2) plt.xlabel(Olympic Year) plt.ylabel(Gold Medals Count) plt.title(Top 10 Countries Gold Medal Trend (1896-2020)) plt.legend(bbox_to_anchor(1.05, 1), locupper left) plt.grid(True, alpha0.3) plt.tight_layout() plt.savefig(olympic_gold_trend.png, dpi300, bbox_inchestight) plt.show()5.2.1 生成毕业设计报告核心图表聚类结果雷达图对KMeans聚类的三类国家计算每类的平均特征值绘制雷达图直观展示差异import numpy as np # 获取聚类结果从4.1节 cluster_summary result_df.groupBy(cluster) \ .agg( avg(avg_gold).alias(avg_gold), avg(avg_growth).alias(avg_growth), avg(avg_total).alias(avg_total) ).toPandas() # 雷达图数据准备 labels [Avg Gold, Avg Growth Rate, Avg Total Medals] stats cluster_summary[[avg_gold, avg_growth, avg_total]].values.T.tolist() # 绘制 angles [n / float(len(labels)) * 2 * np.pi for n in range(len(labels))] stats stats[:1] # 闭合图形 angles angles[:1] ax plt.subplot(111, polarTrue) ax.plot(angles, stats[0], linewidth2, labelEmerging) ax.fill(angles, stats[0], alpha0.25) ax.plot(angles, stats[1], linewidth2, labelTraditional) ax.fill(angles, stats[1], alpha0.25) ax.plot(angles, stats[2], linewidth2, labelDeclining) ax.fill(angles, stats[2], alpha0.25) ax.set_xticks(angles[:-1]) ax.set_xticklabels(labels) ax.legend(locupper right, bbox_to_anchor(1.3, 1.0)) plt.title(Country Competitiveness Clusters Radar Chart) plt.savefig(cluster_radar.png, dpi300, bbox_inchestight)注意雷达图需matplotlib3.7支持极坐标填充。若环境受限改用柱状图对比三类国家的平均特征值代码更简洁且信息量不减。本文还有配套的精品资源点击获取
返回列表