Python+Spark奥运会数据分析与可视化:从数据清洗到ECharts展示
2026/9/17 2:27:52 网站建设 项目流程

简介:基于Python和Spark的奥运会可视化分析系统毕业设计项目,面向计算机专业学生与大数据方向开发者,提供从数据清洗、聚合分析到可视化展示的完整实现方案;项目已在Windows 10和Windows 11环境完成调试,并配有详细部署说明,下载解压后即可运行,也可直接用于课程设计或答辩演示。资源包为ZIP压缩格式,共60个文件,整体体积约1.62MB;包内除核心Python脚本及编译文件外,还包含Spark与Scala处理程序、MySQL数据库脚本、CSV格式数据集、XML与Properties配置文件,以及CSS、JavaScript和图片等前端展示资源,目录层次清晰,便于逐模块研读。项目曾获导师认可,答辩评分达到97分,具备较强的参考价值。预览信息显示,源码中整合了Flask应用、Hadoop与Spark的金牌数据分析模块,并附带ECharts可视化组件和Maven工程配置,方便二次开发与功能扩展。目前已有434人学习下载,对准备毕业设计或希望快速上手大数据可视化实战的读者来说,是一份可以直接借鉴的高质量项目。

1. 用 Python+Spark 做奥运会可视化分析,先想清楚这三件事

很多同学拿到“基于Python+Spark的奥运会可视化分析系统”这个题,第一反应是找一份CSV,画两个ECharts柱状图,然后写一堆Spark代码撑场面。但一个能拿高分的毕业设计,关键是让Spark真正参与数据处理,而不是做一个挂着Spark名字的Excel看板。这个系统要解决的是奥运历史数据从清洗、聚合到前端可视化展示的完整链路:数据量不算大,但用Spark运行出的分区、shuffle、缓存和内存调优过程,恰恰是答辩时最值得讲的部分。我会按平时接这类项目的做法,从环境、清洗、调优到Flask+ECharts,给你一套能复现的代码路径。适合准备毕业设计或面试想讲透Spark的同学。

2. 环境与数据准备:本地模式跑通 PySpark 读取奥运会 CSV 的最小命令

2.1 用 venv + pip 安装 PySpark 并准备好 JDK

PySpark 虽然面向 Python 开发,但底层还是 JVM。所以最先需要处理的是 JDK 版本,而不是直接 pip install。我一般先装 JDK 11,再建虚拟环境安装 PySpark 3.5.x,因为 3.5 对 Python 3.8-3.11 的支持比较稳定,Java 8/11/17 都能跑。如果电脑上已经装了多个 Java 版本,可以用java -version确认默认版本,避免 Spark 启动时报 UnsupportedClassVersionError。

java -version python -m venv venv source venv/bin/activate # Windows 下执行 venv\Scripts\activate pip install pyspark==3.5.1 pandas pyarrow

这里把 pandas 和 pyarrow 一起装上,是为了后面对拍结果、读取 Parquet 时少折腾。安装完成后,不需要单独下载 Spark 发行包,PySpark 自带了 local 模式所需的全部 JAR,对本地开发展够用了。如果以后要连集群,再用spark-submit --master yarn指向集群,环境中需要额外部署 Spark 客户端。

2.2 用 SparkSession 读取奥运会 CSV 并打印 Schema

先确认数据列。公开的奥运历史数据一般包含 Year、City、Sport、Event、Athlete、Country、Medal 这些字段,Medal 为 Gold/Silver/Bronze/NA。读取时用 SparkSession 的 DataFrameReader,指定 header 和 inferSchema 即可,不要直接spark.read.csv不带参数,否则第一行会被当数据读到列名为_c0

from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("Olympic Analysis") \ .master("local[*]") \ .config("spark.sql.shuffle.partitions", "4") \ .getOrCreate() df = spark.read.format("csv") \ .option("header", True) \ .option("inferSchema", True) \ .load("data/olympics.csv") df.printSchema() df.select("Year", "Country", "Medal").limit(5).toPandas()

local[*]表示使用本机所有可用核心,毕业设计演示机不用调太高。spark.sql.shuffle.partitions是聚合和 join 时的分区数,本地 4 个分区比较合适,设置为 200 在本地反而会拖慢小任务。inferSchema 会在读取时扫描一遍文件推断字段类型,如果数据文件格式较规整,可以接受;字段很多时建议改成手工指定 schema。

参数作用本地推荐值
header第一行作为列名True
inferSchema自动推断字段类型True
quote指定字符串引用符号"
multiLine允许字段跨行False

如果读取后字段长度不符合预期,可以用df.describe().show()快速看数值列的最大最小和统计量。

2.3 设置 driver 内存与执行内存,避免演示时进程崩掉

本地模式里,PySpark 进程和 JVM 是同一台机器,driver 和执行器共享内存。若数据量百万级以下,默认内存一般够用;但如果同时打开多个 Spark Shell,容易出现 Java heap space。这时配置两个关键项。

spark = SparkSession.builder \ .master("local[2]") \ .config("spark.driver.memory", "1g") \ .config("spark.executor.memory", "2g") \ .getOrCreate()

spark.driver.memory给 driver 的堆内存,spark.executor.memory是执行器的堆内存。在本地 local 模式下,executor 内存会叠加到同一个 JVM,所以两个值不要同时设过大,否则可能超过本机物理内存。建议在 8G 机器上用 driver 1g / executor 2g,在 16G 机器上可以升到 driver 2g / executor 4g。此外,Spark UI 默认在http://localhost:4040展示任务进度,后面验证阶段还会用到这个界面。

如果后续要处理更大数据,需要调整的就不只是内存,还有spark.executor.coresspark.executor.instances这些集群级参数。本地先把链路跑通,再把同一套代码用spark-submit提交到集群,是更稳的做法。

3. 用 Spark DataFrame 做清洗与聚合:奖牌榜、趋势与选手画像的代码实现

3.1 先做数据质量检查:空值、重复与字段统计

数据清洗不是从 groupBy 开始,而是先看字段质量。olympics.csv这类历史数据常见问题有:Medal 字段为 NA 表示未获奖;Age 列为空;Country 存在前后空格或历史改名;Event 字段带换行符。先用 count 和 countDistinct 看每个字段的有效性。

from pyspark.sql.functions import countDistinct, col, count, when df.select([countDistinct(col(c)).alias(c + "_distinct") for c in df.columns]).show(truncate=False) df.select([count(when(col(c).isNull(), c)).alias(c + "_null") for c in df.columns]).show(truncate=False)

第一句计算每列的非重复值数量,帮助判断哪些字段适合作为维度。第二句统计空值数,注意 DataFrame 里空字符串''不算 null,如果数据里用空字符串表示缺失,还要追加trim(col(c)) == ''条件。用truncate=False是为了让长字符串字段不被截断,避免看到一堆省略号。

3.2 用 withColumn 与 filter 完成空值和异常值处理

我习惯把清洗逻辑写成一个函数,方便在文档里描述和后面复用。清洗步骤这样安排:

from pyspark.sql.functions import trim, upper, col, when, coalesce, lit def clean_olympics(df): return df.filter(col("Year").isNotNull()) \ .withColumn("Country", upper(trim(coalesce(col("Country"), lit("Unknown"))))) \ .withColumn("Medal", when(col("Medal").isin("Gold", "Silver", "Bronze"), col("Medal")).otherwise("NA")) \ .filter(col("Age").isNull() | ((col("Age") >= 6) & (col("Age") <= 100))) df_clean = clean_olympics(df)

每一行的含义:Year 为空直接删掉,因为年份是后续趋势分析的必需维度。Country 先去掉首位空格再转大写,避免 “us” 和 “USA” 被当成两个国家。Medal 只保留三种有效奖牌,其余统一成 NA,否则后续filter(col("Medal") != "NA")会漏掉大小写不同变体。Age 字段对异常范围做过滤,小于 6 或大于 100 的属于录入异常,参与年龄画像统计会拉偏均值。

清洗之后最好做一次df_clean.cache()并执行一个 action(比如df_clean.count()),让后面多次分析复用同一份缓存数据。需要留意:如果数据总量比执行器内存大,缓存策略应该用persist(StorageLevel.MEMORY_AND_DISK),避免 OOM。

3.3 核心分析指标:奖牌榜、历年金牌趋势与年龄分布

下面这组代码是可视化系统里最常用的三个统计结果,也是 spark 数据分析案例里最典型的分组聚合场景。

from pyspark.sql.functions import count, sum, row_number, col from pyspark.sql.window import Window # 奖牌总数 top10 medal_top10 = df_clean.filter(col("Medal") != "NA") \ .groupBy("Country") \ .agg(count("Medal").alias("total_medals")) \ .orderBy(col("total_medals").desc()) \ .limit(10) # 历年金牌趋势:每年每个国家的金牌数 gold_by_year = df_clean.filter(col("Medal") == "Gold") \ .groupBy("Year", "Country") \ .agg(count("Medal").alias("gold_count")) # 每年金牌榜第一名 window_spec = Window.partitionBy("Year").orderBy(col("gold_count").desc()) champion = gold_by_year.withColumn("rank", row_number().over(window_spec)) \ .filter(col("rank") == 1) \ .select("Year", "Country", "gold_count")

window_spec按年分区,在分区内按金牌数降序排名,每年只保留第 1 名。这个查询在纯 DataFrame 写法里不需要显式调整 RDD partition,Spark 会自动根据 shuffle.partitions 来分布。如果想给前端同时提供 top 榜和冠军历史两张图,把结果分别写出即可。

3.4 导出分析结果:Parquet、CSV 与 toPandas 的取舍

可视化层一般不需要直接连接 Spark,数据量不大时可以先导出为文件再交给 Web 服务。常见的做法是聚合结果写 Parquet,因为列式存储压缩率高,而且 Spark 或 pandas 读起来都快。小型排行榜可以直接 toPandas 转 JSON 返回,但要注意 toPandas 会在 driver 端拉取全量数据,大数据集下必然内存溢出。

medal_top10.write.mode("overwrite").parquet("output/medals_top10.parquet") gold_by_year.write.mode("overwrite").parquet("output/gold_by_year.parquet") df_clean.groupBy("Age").count().write.mode("overwrite").csv("output/age_dist.csv", header=True)

写输出用mode("overwrite")是让重复运行不会因为输出目录存在而失败。如果需要给前端吐 JSON,可以用toJSON或者收集到 Python 端。这里给出一个比较基准:

输出格式适用场景注意点
Parquet后续 Spark/pandas 再读取需要 pyarrow,非文本格式
CSV给 Excel 或文档截图字段含逗号时要加 quote
JSON LinesWeb API 接口消费一条记录一行,Spark 的to_json生成

把分析结果落到文件后,第 5 章再做可视化时只需读结果文件,不需要在每次请求时启动 SparkSession。这个边界一定提前告诉队友或写在文档里,否则有人会把 Spark 调优参数和 Flask 请求混在一起,接口响应时间会很难看。

4. 性能调优与数据倾斜:让 Spark 分析奥运会数据更稳的三类参数

4.1 从 Spark UI 和执行计划定位性能瓶颈

页面卡住是常态,但卡在哪个阶段要弄明白。Spark 提供了 Web UI,在 local 模式下跑任务时访问http://localhost:4040,能看到每个 Job、Stage、Executor 的耗时和 shuffle 读写量。毕业设计答辩时如果能指着 UI 讲“这个 Stage 的 Shuffle Write 是 XXX,说明我调整了并行度”,比说“我用了分布式计算”有说服力得多。

champion.explain("formatted")

explain输出逻辑计划和物理计划,重点看 Exchange 节点。如果一个 groupBy 只有几个 key,但 Exchange 输出特别大,很可能发生了数据倾斜。拿奥运会数据举例,按 Country 分组时美国和历年东道主记录条数明显多于小国,少数分区会拖慢整个 Stage。

4.2 分区与并行度:repartition、coalesce 和 shuffle.partitions

很多同学一上来就调spark.executor.memory,其实大多数慢在分区数不合适。分区太少,每区数据量不均衡;分区太多,调度和序列化开销反而变大。我一般这样判断:

# 查询重分区为 8 个分区,按年份保证同一国家在相同分区 df_repartitioned = df_clean.repartition(8, "Year") # 减少到 4 个分区,coalesce 不做全量 shuffle df_reduced = df_repartitioned.coalesce(4)

repartition(n, col)会触发全节点 shuffle,利用 key 将相同 key 的数据放到同一任务,避免 join 时跨节点传输。coalesce(n)只把现有分区合并,不产生 full shuffle,适合在数据量变小的阶段使用。下列参数直接影响聚合类任务的行为:

参数名作用本地推荐值
spark.sql.shuffle.partitions聚合/join 的分区数4-8
spark.default.parallelismRDD 默认并行度与核数相关
spark.sql.adaptive.enabled动态调整分区true
spark.sql.adaptive.coalescePartitions.enabled自动合并小分区true

4.3 用 DataFrame 表达式替代 Python UDF,减少序列化开销

Python UDF 是 PySpark 的“蜜糖陷阱”:写法直观,但每行数据都要在 JVM 和 Python 进程间序列化,性能差一个数量级。奥运会数据量小可能看不出差异,但如果按 spark 集群搭建后的真实数据规模,UDF 会成为瓶颈。优先用内置函数。

# 不推荐:Python 自定义函数处理空值和大小写 from pyspark.sql.udf import udf from pyspark.sql.types import StringType @udf(StringType()) def clean_country(c): return c.strip().upper() if c else "Unknown" # 推荐:直接用 trim/upper/coalesce 完成同样逻辑 from pyspark.sql.functions import trim, upper, coalesce, lit df_clean = df.withColumn("Country", upper(trim(coalesce(col("Country"), lit("Unknown")))))

如果业务逻辑实在复杂,必须用 Python 函数,至少选择 pandas UDF(@pandas_udf),利用 Arrow 批量序列化,避免逐行转换。这是 5 年经验的人也会在代码评审里给的建议。

4.4 缓存与广播变量:什么时候该用,什么时候别用

cache()使用成本很低,但乱用也会让内存管理变差。我的一般规则是:同一个 DataFrame 会被后续 3 个以上作业反复读取才缓存;中间结果在过滤到很小维度后,可以用广播变量参与 join,让每个 executor 保留一份小表副本,避免大 shuffle。

from pyspark.sql.functions import broadcast medal_dict = df_clean.select("Medal").distinct().collect() broadcast_medal = spark.sparkContext.broadcast([row.Medal for row in medal_dict]) df_clean = df_clean.withColumn("medal_flag", col("Medal").isin(*broadcast_medal.value))

这里把奖牌列表广播到各 executor,洗数据时就不需要全局分布式查字典。用broadcast(df_dim)进行 join,是数据倾斜场景下更快的手段。缓存则要注意:同一个分析里聚合之后再读一次结果文件,而不是把清洗后的全量数据缓存后反复跑 groupBy,后者会让缓存占用大于实际收益。

5. Flask + ECharts 可视化:把 Spark 计算结果变成奖牌榜与趋势图

5.1 可视化层只读结果,不连 Spark

在前面几步的架构中,Spark 只负责离线计算,Web 服务用 Flask 读取第 3 章导出的 Parquet/CSV 文件并生成接口。这个设计有两个直接收益:第一,每次刷新图表不会重新启动 SparkContext,避免 5-10 秒的初始化延迟;第二,可视化服务和数据处理可以分开调试,一个人负责算数、一个人负责画图,互不阻塞。系统的目录结构可以是这样:

. ├── app.py ├── templates/ │ └── index.html ├── static/ │ └── js/ │ └── dashboard.js ├── output/ │ ├── medals_top10.parquet │ ├── gold_by_year.parquet │ └── age_dist.csv └── data/ └── olympics.csv

这种分层还被很多线上项目沿用:数据层、计算层、服务层、展示层各自独立,适合写进文档架构图。

5.2 用 Flask 暴露奖牌榜和趋势 JSON 接口

后端代码很短。读 Parquet 需要 pandas 和 pyarrow,前面已经安装;如果结果文件是 CSV,直接用 pandas.read_csv 也可以。

from flask import Flask, jsonify, render_template import pandas as pd app = Flask(__name__) def read_medal_top10(): df = pd.read_parquet("output/medals_top10.parquet") return {"categories": df["Country"].tolist(), "values": df["total_medals"].astype(int).tolist()} @app.route("/") def index(): return render_template("index.html") @app.route("/api/medals") def medals_api(): return jsonify(read_medal_top10()) @app.route("/api/trend") def trend_api(): df = pd.read_parquet("output/gold_by_year.parquet") us = df[df["Country"] == "USA"].sort_values("Year") return jsonify({"year": us["Year"].tolist(), "gold_count": us["gold_count"].tolist()}) if __name__ == "__main__": app.run(host="0.0.0.0", port=5000, debug=True)

/api/trend先筛选 USA 再按年份排序,生成折线图需要的一维数组。注意debug=True仅用于开发,演示时开着调试模式可能被误触发 reloader。生产环境建议把 debug 关掉,并用app.run(host="0.0.0.0", port=5000)让局域网内其他设备也能访问页面。

接口路径方法返回字段前端用途
/api/medalsGETcategories, values柱状图
/api/trendGETyear, gold_count折线图

5.3 ECharts 从接口取数,绘制 top10 柱状图

前端用 ECharts 的 init 和 fetch,逻辑清晰。下面是templates/index.html的最小片段:

<!DOCTYPE html> <html lang="zh"> <head> <meta charset="UTF-8"> <title>奥运会可视化分析</title> <script src="https://cdn.jsdelivr.net/npm/echarts@5/dist/echarts.min.js"></script> </head> <body> <div id="medal_chart" style="width: 900px; height: 500px;"></div> <script> fetch('/api/medals') .then(res => res.json()) .then(data => { const chart = echarts.init(document.getElementById('medal_chart')); chart.setOption({ title: { text: '国家奖牌总数 Top10' }, xAxis: { type: 'category', data: data.categories }, yAxis: { type: 'value' }, series: [{ type: 'bar', data: data.values }] }); }); </script> </body> </html>

这段代码在加载页面时向/api/medals发请求,拿到 categories 和 values 后渲染柱状图。毕业设计阶段不需要引入复杂状态管理,保持一个接口对应一个图表即可。

注意:不要在浏览器里直接双击打开 HTML,ECharts 请求 file:// 会触发跨域报错,正确方式是先启动 Flask,再访问 http://localhost:5000/。

5.4 用折线图展示历年冠军走势,补齐“分析系统”的维度

光有排行不算可视化分析系统,再加一张历年金牌趋势图更能体现分析能力。在dashboard.js里复用同样的 fetch 逻辑,请求/api/trend并初始化折线图:

fetch('/api/trend') .then(res => res.json()) .then(data => { const chart = echarts.init(document.getElementById('trend_chart')); chart.setOption({ title: { text: 'USA 历年金牌数趋势' }, xAxis: { type: 'category', data: data.year }, yAxis: { type: 'value' }, series: [{ type: 'line', data: data.gold_count, areaStyle: {} }] }); });

为了让趋势更有意义,除了 USA 还可以在前端条件筛选多个国家。后端可加一个参数/api/trend?country=CHN,Flask 用request.args.get("country", "USA")获取,Spark 结果文件已在内存中,筛选成本非常低。到这里,Pandas 读 Parquet、Flask 发 JSON、ECharts 成图,链路已经完整。

6. 验证、Spark UI 与毕业设计文档收尾技巧

6.1 用 explain 和 Spark UI 验证计算正确性

前面用过 explain,验证阶段更要用它。对拍方法最直接,把 Spark 聚合结果和 Pandas 对原始 CSV 的聚合结果 diff 一下。

expected = pdf[pdf.Medal.notna()].groupby("Country").size().nlargest(10) actual = medal_top10.toPandas().set_index("Country")["total_medals"] diff = abs(expected - actual) assert diff.max() == 0

这里的前提是 Spark 清洗和 Pandas 处理逻辑一致,比如空值过滤、奖牌判定都要对齐。如果 diff 不为零,先回头检查是否存在大小写和空格差异。作业跑完后在 Spark UI 的 SQL/Job 标签页看各 Stage 的 Shuffle Read 和 Write,中间阶段异常时会把瓶颈暴露出来。

6.2 处理数据倾斜的 3 个快速手段

如果是按国家分组,大国记录多到让某个任务耗时独大,有三个临时手段:第一,repartition("Country", 8)在分组前把大 key 分散到多分区;第二,过滤掉明显无用的数据,比如只保留 1980 年以后的赛事;第三,对真正热点 key 加两阶段聚合的盐值,先按Country + salt聚合一次,再去盐按Country聚合第二次。第一种最省事,第三种适合你认为这个倾斜点会写在论文里的场景,不建议对所有 key 都做。

6.3 把调试过程整理成文档和答辩截图

使用文档里除了写“如何启动 Flask、如何修改端口”,还要放 3 张图:Spark UI 的 Job 列表、ECharts 图表页面、以及结果目录的 Parquet/CSV 文件。截图时注意用浏览器无痕窗口打开页面,避免地址栏出现 file:// 路径。图注格式建议写清“图 3-2 2024 年奥运会奖牌榜 Top10 柱状图”,不要用中文逗号分隔编号和标题。文档里再补一张 Excel 格式的参数对照表,把 spark.executor.memory、spark.sql.shuffle.partitions 和实际机器配置放到同一行,答辩老师问调优时直接指给它看。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询