ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

基于奥运会奖牌数据的Spark分析与可视化实战

2026/9/12 10:46:48 拓冰建站 浏览量
基于奥运会奖牌数据的Spark分析与可视化实战 简介面向大数据、人工智能及相关专业毕业设计需求这份项目完整实现了基于Hadoop与Spark的奥运会奖牌变化分析案例从奖牌历史数据预处理、统计聚合、趋势可视化到基于Flask的Web结果展示并配齐MySQL数据库脚本、Scala/Java工程配置和数据集文件可作为毕设、课设或项目初期演示直接使用也方便二次开发。资源共60个文件以Python源码py、Scala源码、Java配置xml/properties、SQL脚本与CSV数据集为主另有截图、说明文档和运行所需依赖清单压缩包仅1.62MB目录结构清晰便于按模块查阅。目前已有93人学习下载适合具备一定大数据基础、希望快速搭建完整分析链路的学生或开发者。配套材料覆盖设计文档、完整源码、数据库文件、授权说明等能帮助使用者省去大部分环境搭建与数据准备时间集中理解Spark分析逻辑、Hadoop生态协作方式以及整体系统架构对毕业答辩和实际项目落地都有直接参考价值。1. 从奥运会奖牌数据看 Hadoop/Spark 选型一份跨越 1896 到 2016 年的奥运会奖牌数据集文件只有几 MB很多人觉得没必要上 Hadoop。但真正跑过数据的人清楚瓶颈不在存储,而在清洗和维度聚合的复杂度运动员改国籍、项目更名、联合组队这些脏数据叠在一起需要反复迭代计算。Spark 的价值是把清洗、聚合、回写数据库这条链路统一成一份可复现的 DataFrame 代码用完即走不必像传统数仓那样先建模再灌数据。这套 Spark 数据分析案例面向大数据专业学生和刚接触 Spark 的开发者用奥运会这个容易理解的业务场景把 Hadoop 生态下的数据接入、Spark 计算、MySQL 存储、Flask 接口串成完整链路。拿到源码后替换数据源和表结构就能迁移到其他赛事分析场景。2. 数据接入与存储设计Olympic 数据集清洗与 MySQL 表结构拿到压缩包先解压olympicSummer 2是原始 CSV 数据集book.sql是 MySQL 初始化脚本app.py是 Flask 入口。建议按「先建库 → 再清数据 → 后跑 Spark」的顺序操作顺序反了容易出现 MySQL 表还没建好、Spark 任务却已经在写库的尴尬情况。2.1 源数据集字段解析与常见脏数据olympicSummer 2目录下的 CSV 是夏季奥运会历史奖牌明细核心字段如下字段名示例说明CityRio举办城市Year2016举办年份SportAquatics运动大项DisciplineSwimming分项EventMens 100m Freestyle小项AthleteMichael Phelps运动员姓名GenderMen性别CountryUSA国家/地区代码MedalGold奖牌类型导入 MySQL 之前先做两步检查。第一步用head -5 olympicSummer\ 2.csv看表头和分隔符确认是标准逗号分隔而不是制表符第二步统计空值和异常国家代码。这个数据集最常见的坑是早期奥运会没有女子项目Gender字段为空部分运动员姓名含重音符号不带 UTF-8 编码导入会出现乱码。提示CSV 里的 Year 字段部分行混入空字符串Spark 的inferSchema可能把整列推断成 string清洗时要对年份做一次类型校验。清洗时建议用 Spark 读一次原始数据把明显异常的行过滤掉输出一份清洗报告而不是直接改原文件。这样换数据集时只需调整过滤条件不需要重写导入逻辑。先用wc -l看总行数再用du -h看文件大小典型夏季奥运会奖牌明细在三万行以内、单文件不到 5 MB。这个量级用 HDFS 略重但课程设计的重点不是性能而是完整走一遍大数据处理流程把 CSV 放到 HDFS、用 Spark 读 HDFS 路径、结果写 MySQL。2.2 MySQL 表结构设计与建表 SQL存储层我设计了四张表olympic_medals存奖牌明细country_year_medals存国家年度聚合结果country_medal_trend存趋势计算后的环比数据dim_country存国家代码映射。明细表保留全量原始字段聚合表只保留分析维度。CREATE DATABASE IF NOT EXISTS olympic DEFAULT CHARACTER SET utf8mb4; CREATE TABLE olympic_medals ( id INT PRIMARY KEY AUTO_INCREMENT, city VARCHAR(64), year INT, sport VARCHAR(64), discipline VARCHAR(64), event_name VARCHAR(128), athlete VARCHAR(128), gender VARCHAR(16), country VARCHAR(8), medal VARCHAR(16), KEY idx_year_country (year, country), KEY idx_country (country) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; CREATE TABLE country_year_medals ( id INT PRIMARY KEY AUTO_INCREMENT, country VARCHAR(8) NOT NULL, year INT NOT NULL, gold INT DEFAULT 0, silver INT DEFAULT 0, bronze INT DEFAULT 0, total INT DEFAULT 0, UNIQUE KEY uk_country_year (country, year) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;参数上有几个设计点。country_year_medals加联合唯一键uk_country_year是为了让 Spark 回写聚合结果时走「存在则更新、不存在则插入」的幂等逻辑避免重复任务产生重复数据。olympic_medals上的idx_year_country索引服务 Flask 接口按年份和国家筛选避免全表扫描。字符集用utf8mb4而不是utf8因为utf8mb4能覆盖四字节字符运动员姓名中的生僻字和中文注释都能正常存储。2.3 数据导入Spark 读 CSV 与 MySQL 写入清洗和导入的核心逻辑用 Spark DataFrame 实现不写 MapReduce 是因为这步要处理字段裁剪、类型校验和编码问题DataFrame 的声明式 API 比 MapReduce 的 map/reduce 原语直观得多。源码里是 Java 版本这里给出等价的 Scala 片段逻辑一致import org.apache.spark.sql.{SparkSession, SaveMode} val spark SparkSession.builder() .appName(OlympicMedalImport) .master(local[*]) // 本地模式调试集群提交改为 yarn .getOrCreate() val df spark.read .option(header, true) .option(inferSchema, true) .option(encoding, UTF-8) .csv(hdfs:///data/olympicSummer_2.csv) val cleanDf df.filter( $Country.isNotNull $Medal.isNotNull $Year.cast(int).isNotNull )inferSchema会自动推断列类型但年份列混入空串时推断可能失败所以 filter 里用cast(int)再做一次校验。加载后写 MySQLcleanDf.select( $City, $Year, $Sport, $Discipline, $Event.alias(event_name), $Athlete, $Gender, $Country, $Medal ).write .mode(SaveMode.Append) .option(batchsize, 1000) .jdbc(jdbc:mysql://localhost:3306/olympic, olympic_medals, new java.util.Properties() {{ setProperty(user, root) setProperty(password, your_password) setProperty(driver, com.mysql.cj.jdbc.Driver) }})batchsize控制每批写入行数MySQL 8.0 下设 1000 比默认快 30% 左右。Append模式保证重复导入只追加不覆盖。导入完成后执行SELECT COUNT(*) FROM olympic_medals和源文件行数比对两个数字必须一致不一致就回头查过滤条件。3. Spark 核心计算奖牌变化趋势与国家维度聚合3.1 为什么用 DataFrame 而不是 RDD如果这个案例用纯 Hadoop MapReduce 写需要两个 Job第一个统计每个国家每年的奖牌数量第二个按年份排序计算变化。两个 Job 之间要把中间结果落到 HDFS代码量翻倍。Spark 把这两步压缩到一个 DAG 里窗口函数和聚合在同一个执行计划内完成这是推荐用 Spark 做主计算的原因。具体到代码选型DataFrame 比 RDD 更合适它有 Catalyst 优化器做谓词下推和列剪枝写起来也更接近 SQL 语义。如果需要更细粒度的类型安全可以换 Dataset但这个案例字段结构简单DataFrame 完全够用。3.2 计算各国年度奖牌数先算每个国家每年的金银铜和总数Scalaimport org.apache.spark.sql.functions._ // 条件计数效果等同 CASE WHEN val medalAgg cleanDf.groupBy(Country, Year) .agg( sum(when($Medal Gold, 1).otherwise(0)).alias(gold), sum(when($Medal Silver, 1).otherwise(0)).alias(silver), sum(when($Medal Bronze, 1).otherwise(0)).alias(bronze), count($Medal).alias(total) )when().otherwise()做条件计数效果等同于 SQL 里的CASE WHEN。输出字段如下输出字段含义计算方式gold金牌数Medal Gold 条件计数silver银牌数Medal Silver 条件计数bronze铜牌数Medal Bronze 条件计数total奖牌总数分组内全部记录计数如果熟悉 SQL可以写spark.sql(SELECT Country, Year, SUM(CASE WHEN MedalGold THEN 1 ELSE 0 END) ...)两种写法经过 Catalyst 优化后执行计划一致。3.3 计算奖牌变化率lag 窗口函数接着算每个国家相对上一届奥运会的奖牌变化量。这里有个坑奥运会每四年一届但一战、二战导致停办相邻两条记录的年份差不一定等于 4。窗口函数必须按Country分区、按Year排序然后用 lag 取前一行再用年份差过滤import org.apache.spark.sql.expressions.Window val winSpec Window.partitionBy(Country).orderBy(Year) val trendDf medalAgg .withColumn(prev_total, lag(total, 1).over(winSpec)) .withColumn(prev_year, lag(Year, 1).over(winSpec)) .withColumn(year_gap, $Year - $prev_year) .filter($year_gap 4 || $year_gap.isNull) .withColumn(change, $total - coalesce($prev_total, lit(0)))lag(total, 1)取分区内前一行over(winSpec)定义窗口边界。过滤year_gap 4只保留正常奥运周期停办造成的非正常间隔被排除。coalesce把第一年的prev_total空值补成 0让首年的 change 显示为 total而不是 null。提示同一个 Event 下出现多个 Bronze 记录是并列铜牌属于正常业务规则清洗时不要对整行做 distinct()。3.4 结果写入 MySQL 与校验趋势结果写入country_medal_trend表代码与上一章 JDBC 写入一致区别是写模式用SaveMode.Overwrite防止多次任务执行后趋势数据越积越多。落库后用下面 SQL 验证中美两国的金牌变化趋势SELECT country, year, gold, total FROM country_medal_trend WHERE country IN (CHN, USA) AND year 2000 ORDER BY country, year;对照国际奥委会官网奖牌榜中国在 2008 年北京奥运会金牌数达到峰值2012 年回落2016 年继续下降。如果跑出来的数据和历史事实对得上说明清洗和聚合逻辑基本正确。需要注意数据集里部分国家代码是旧的例如苏联解体前的 URS这些代码在 1992 年后不应继续出现查询时用dim_country表过滤已废弃代码。4. Flask API 层设计与可视化联动从 MySQL 到浏览器4.1 API 接口设计与 Flask 路由Spark 任务是离线批处理Flask 服务是实时查询两者通过 MySQL 解耦。聚合结果放在 MySQL 而不是 HBase是因为奖牌数据量小、查询模式固定用关系型数据库的 SQL 就能覆盖答辩时演示也更直观。接口设计为三个接口路径方法参数返回内容/api/trendGETcountry, start_year, end_year指定国家奖牌变化趋势/api/rankGETyear, medal_type某一年奖牌榜排名/api/countriesGET无国家代码列表前端下拉框用app.py的关键片段from flask import Flask, jsonify, request import pymysql app Flask(__name__) db_config { host: localhost, port: 3306, user: root, password: your_password, database: olympic, charset: utf8mb4, cursorclass: pymysql.cursors.DictCursor, } app.route(/api/trend, methods[GET]) def get_trend(): # 从查询字符串读取参数带默认值方便直接访问 country request.args.get(country, CHN) start request.args.get(start_year, 1996, typeint) end request.args.get(end_year, 2016, typeint) sql SELECT year, gold, silver, bronze, total FROM country_year_medals WHERE country %s AND year BETWEEN %s AND %s ORDER BY year conn pymysql.connect(**db_config) try: with conn.cursor() as cur: cur.execute(sql, (country, start, end)) # 参数化查询防注入 rows cur.fetchall() finally: conn.close() return jsonify({country: country, data: rows})这里最重要的点是 SQL 参数化%s占位符由 pymysql 转义避免countryUSA OR 11这类注入。DictCursor让查询结果直接是字典列表省去手动组装 JSON。如果查询结果超过几万条加LIMIT或用stream_with_context流式返回但奥运奖牌场景不会到这一步。接口最好再包一层统一响应结构例如{code: 0, message: ok, data: rows}方便前端做统一错误处理也便于答辩时快速定位是后端异常还是前端渲染问题。4.2 前端 ECharts 折线图联动前端在static/index.html里用 ECharts 渲染折线图。页面加载时先请求/api/countries拉国家列表用户选中某国后发起/api/trend请求动态更新 seriesasync function loadTrend(country) { const resp await fetch(/api/trend?country${country}start_year1996end_year2016); const json await resp.json(); // 动态更新折线图序列 chart.setOption({ xAxis: { type: category, data: json.data.map(d d.year) }, series: [{ name: 总奖牌数, type: line, data: json.data.map(d d.total), smooth: true }] }); }fetch的 URL 里直接拼接${country}虽然后端已经参数化查询但前端最好加一层encodeURIComponent(country)避免国家代码含特殊字符时请求路径被截断。setOption重复调用做的是增量更新不会重置已经配置好的坐标轴样式。如果还想展示金牌、银牌、铜牌各自的走势可以在 series 数组里继续追加对象data 字段分别取d.gold、d.silver、d.bronze。4.3 跨域与联调注意事项Flask 默认监听 127.0.0.1:5000前端如果用 Vite 或 Webpack dev server 跑在 5173 端口浏览器会触发 CORS 拦截。两种解决方式后端装flask-cors执行CORS(app)或者把前端构建产物放到 Flask 的static目录由同一端口服务。毕业设计演示选后者更稳答辩现场不会因为跨域问题翻车。接口写完先用 curl 做冒烟测试例如curl -s http://127.0.0.1:5000/api/trend?countryUSA确认返回 JSON 结构和字段名与前端预期一致再打开页面联调能省掉不少调试时间。5. 参数调优与验收spark-submit 提交和结果一致性校验5.1 从本地模式切换到 YARN 集群提交开发时 IDE 里直接跑master(local[*])没问题Hadoop 集群搭建完成后正式任务要走 YARN 提交。将代码打成olympic-spark.jar后执行spark-submit \ --master yarn \ --deploy-mode cluster \ --name olympic-medal-analysis \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 4 \ --conf spark.sql.shuffle.partitions20 \ --class com.example.OlympicAnalysis \ olympic-spark.jar--driver-memory给 Driver 存放查询计划和执行上下文使用本案例结果直接写 MySQL不在 Driver 端 collect 大结果集2g 足够。--executor-memory 4g在这个数据量下绰绰有余。真正影响效率的是spark.sql.shuffle.partitions默认 200 个分区奖牌数据量只有几千行200 个分区会产生大量无意义的小任务改成 20 让每个分区数据更均衡配合coalesce还能减少写 MySQL 的连接数。5.2 结果一致性校验方法任务跑完不要急着截图先在 MySQL 里验证汇总结果和明细对得上再和官网奖牌榜比对。下面这条 SQL 把聚合表和明细表的数量做交叉检查是最快的对账方法-- 校验聚合表与明细表一致性 SELECT cym.country, cym.year, cym.total, (SELECT COUNT(*) FROM olympic_medals om WHERE om.country cym.country AND om.year cym.year) AS detail_cnt FROM country_year_medals cym LIMIT 10;如果total和detail_cnt对不上问题大概率出在清洗阶段。注意并列铜牌是正常业务规则同一 Event 下多个 Bronze 记录不能整行去重。2016 年里约奥运会美国队 46 金、英国 27 金排名对不上时优先检查奖牌类型字段大小写gold还是Gold。5.3 演示提速缓存与广播变量演示时要来回切换国家和年份二次查询全部落在 MySQL 上Spark 只启动一次这是符合预期的。如果代码里对同一个 DataFrame 多次复用加 cache 能避免重复触发 groupByval cachedAgg medalAgg.cache() val trend cachedAgg.withColumn(prev_total, lag(total, 1).over(winSpec)) val rank cachedAgg.filter($Year 2016).orderBy($gold.desc)cache()返回的是标记了缓存标志的新 DataFrame 引用后面的动作必须从它出发才能命中缓存直接引用原变量不会生效。另外dim_country这种不到 300 行的小表join 时用broadcast(dimCountryDf)包一层Spark 会把这个小表复制到每个 Executor 的本地内存省掉 shuffle 阶段的网络传输。本文还有配套的精品资源点击获取