ARTICLE DETAIL

资讯详情

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

Python+Spark奥运会数据分析与可视化:从数据清洗到ECharts展示

Python+Spark奥运会数据分析与可视化:从数据清洗到ECharts展示 简介基于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. 用 PythonSpark 做奥运会可视化分析先想清楚这三件事很多同学拿到“基于PythonSpark的奥运会可视化分析系统”这个题第一反应是找一份CSV画两个ECharts柱状图然后写一堆Spark代码撑场面。但一个能拿高分的毕业设计关键是让Spark真正参与数据处理而不是做一个挂着Spark名字的Excel看板。这个系统要解决的是奥运历史数据从清洗、聚合到前端可视化展示的完整链路数据量不算大但用Spark运行出的分区、shuffle、缓存和内存调优过程恰恰是答辩时最值得讲的部分。我会按平时接这类项目的做法从环境、清洗、调优到FlaskECharts给你一套能复现的代码路径。适合准备毕业设计或面试想讲透Spark的同学。2. 环境与数据准备本地模式跑通 PySpark 读取奥运会 CSV 的最小命令2.1 用 venv pip 安装 PySpark 并准备好 JDKPySpark 虽然面向 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 pyspark3.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第一行作为列名TrueinferSchema自动推断字段类型Truequote指定字符串引用符号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.cores、spark.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(truncateFalse) df.select([count(when(col(c).isNull(), c)).alias(c _null) for c in df.columns]).show(truncateFalse)第一句计算每列的非重复值数量帮助判断哪些字段适合作为维度。第二句统计空值数注意 DataFrame 里空字符串不算 null如果数据里用空字符串表示缺失还要追加trim(col(c)) 条件。用truncateFalse是为了让长字符串字段不被截断避免看到一堆省略号。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 partitionSpark 会自动根据 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, headerTrue)写输出用mode(overwrite)是让重复运行不会因为输出目录存在而失败。如果需要给前端吐 JSON可以用toJSON或者收集到 Python 端。这里给出一个比较基准输出格式适用场景注意点Parquet后续 Spark/pandas 再读取需要 pyarrow非文本格式CSV给 Excel 或文档截图字段含逗号时要加 quoteJSON 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-8spark.default.parallelismRDD 默认并行度与核数相关spark.sql.adaptive.enabled动态调整分区truespark.sql.adaptive.coalescePartitions.enabled自动合并小分区true4.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 UDFpandas_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(host0.0.0.0, port5000, debugTrue)/api/trend先筛选 USA 再按年份排序生成折线图需要的一维数组。注意debugTrue仅用于开发演示时开着调试模式可能被误触发 reloader。生产环境建议把 debug 关掉并用app.run(host0.0.0.0, port5000)让局域网内其他设备也能访问页面。接口路径方法返回字段前端用途/api/medalsGETcategories, values柱状图/api/trendGETyear, gold_count折线图5.3 ECharts 从接口取数绘制 top10 柱状图前端用 ECharts 的 init 和 fetch逻辑清晰。下面是templates/index.html的最小片段!DOCTYPE html html langzh head meta charsetUTF-8 title奥运会可视化分析/title script srchttps://cdn.jsdelivr.net/npm/echarts5/dist/echarts.min.js/script /head body div idmedal_chart stylewidth: 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 后渲染柱状图。毕业设计阶段不需要引入复杂状态管理保持一个接口对应一个图表即可。注意不要在浏览器里直接双击打开 HTMLECharts 请求 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?countryCHNFlask 用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 和实际机器配置放到同一行答辩老师问调优时直接指给它看。本文还有配套的精品资源点击获取
返回列表