Python爬虫+Spark+Flask:构建农产品价格可视化分析系统
简介基于 Python 爬虫 Flask Spark 的农产品可视化分析系统是面向大数据课程设计的高分项目源码本地编译可运行评审分达 98 分适合毕业设计、期末大作业或课程设计参考。压缩包共 475 个文件涵盖 34 个 Python 后端脚本、17 个 Vue 组件、25 个 JavaScript 前端图表配置以及大量 JSON 数据、Shell 部署脚本和使用说明文档压缩后约 3MB目录结构清晰便于按模块拆解学习。项目完整覆盖 Scrapy 爬虫抓取、Spark 数据处理、Flask 接口服务和可视化展示链路前端包含城市地图、折线图等图表模块可帮助读者快速理解农产品数据从采集、清洗、分析到展示的整套流程并据此替换数据集或扩展功能实现自己的课程设计。目前已有 266 人下载学习适合具备一定 Python 基础、希望高效搭建数据分析可视化项目的中高级学习者。1. 农产品可视化分析的系统定位爬虫、Flask与Spark如何分工做一个农产品可视化分析系统最容易犯的错是把三种技术混在一起用让 Flask 去调度爬虫让 Spark 在 Web 请求里跑计算最后前端拿到的页面卡到超时。真实课程设计里这套系统的价值不在于某个环节多复杂而在于把「采集、计算、展示」拆成清晰的三段——Python 爬虫负责把农产品价格、产地、上市量数据从公开网页抓下来Spark 负责做跨市场、跨品类的聚合统计Flask 只负责把统计结果变成浏览器能看的接口和页面。你要交付的不只是源码还有一条能讲清楚的数据流水线。适合做这个题目的人有两种一是大数据方向的学生需要把 Spark 写进课程设计里体现「分布式计算」能力二是想快速搭一个完整数据链路的工程师用熟悉的技术验证从数据源到可视化大屏的全流程。下面按「爬虫采集 → Spark 分析 → Flask 展示」的顺序把每一步的命令、参数和坑位拆开讲。2. 用Python爬虫把农产品市场数据采下来并清洗入库2.1 选数据源与字段设计先定表结构再写爬虫爬虫不是上来就写 requests而是在写之前想清楚要哪些字段。农产品分析常用的公开数据源有新发地价格行情、地方农业信息网的每日报价、惠农网价格数据等。课程设计一般选一个结构稳定的静态 HTML 页面或 JSON 接口避免选需要登录或强验证的站点否则后期调试成本会吃掉大量时间。建议字段设计如下字段名类型说明product_namestring农产品名称如黄瓜、西红柿marketstring市场名称specstring规格等级如「一级」「统货」min_pricefloat最低价元/斤max_pricefloat最高价元/斤avg_pricefloat平均价元/斤unitstring计价单位report_datestring报价日期格式 yyyy-MM-dd字段定好之后爬虫代码的结构就很清晰了。我的做法是定义ProductItem字典解析完一个条目就扔进列表最后统一批量写入 CSV。下面给一个最小可用的爬虫模板import requests from bs4 import BeautifulSoup import csv import time HEADERS { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36, Referer: http://www.example.com } def fetch_page(url): resp requests.get(url, headersHEADERS, timeout10) resp.raise_for_status() resp.encoding resp.apparent_encoding return BeautifulSoup(resp.text, lxml) def parse_products(html): rows html.select(table.price-list tr) products [] for row in rows[1:]: cells [td.get_text(stripTrue) for td in row.find_all(td)] if len(cells) 6: continue products.append({ product_name: cells[0], market: cells[1], spec: cells[2], min_price: float(cells[3]), max_price: float(cells[4]), avg_price: float(cells[5]), unit: 斤, report_date: time.strftime(%Y-%m-%d) }) return products def save_to_csv(products, filename): with open(filename, a, newline, encodingutf-8-sig) as f: writer csv.DictWriter(f, fieldnamesproducts[0].keys()) if f.tell() 0: writer.writeheader() writer.writerows(products) if __name__ __main__: page fetch_page(http://www.example.com/price) data parse_products(page) if data: save_to_csv(data, products.csv) time.sleep(2)代码的逻辑是先请求页面拿到 HTML再交给 BeautifulSoup 按 CSS 选择器定位价格表格逐行解析出产品字段最后追加写入 CSV。注意encoding resp.apparent_encoding这行用于兼容 GBK 编码的老式农业信息网站遇到中文乱码时优先检查这一项。time.sleep(2)是给目标站的礼貌性限速也避免高频请求触发反爬。2.2 反爬、重试和断点续采爬虫能过夜的三个前提课程设计的数据量通常不大但你要演示的是「能持续采一段时间」所以爬虫的鲁棒性比单次抓取量更重要。常见做法是用requests加重试机制把超时和 5xx 状态码自动重试 3 次页面结构变动时抛出明确异常并记录日志而不是让程序直接崩掉。这里用retry装饰器就行import logging from time import sleep from functools import wraps logging.basicConfig(levellogging.INFO, format%(asctime)s %(levelname)s %(message)s) def retry(times3, delay2): def decorator(func): wraps(func) def wrapper(*args, **kwargs): for i in range(times): try: return func(*args, **kwargs) except Exception as e: logging.warning(第 %s 次请求失败%s, i 1, e) sleep(delay) return None return wrapper return decorator retry(times3, delay2) def fetch_page_safe(url): return fetch_page(url)断点续采的核心是把抓过的日期、分页索引记入本地文件。比如把「当天已经采过哪些市场」写进collected.log下次启动先读日志跳过已完成页面。这样即使爬虫中断也不会从第一条重新抓。2.3 数据清洗别把脏数据直接喂给 Spark爬下来的数据存在三类问题字符串有换行符、均价缺失、同一天同一市场重复记录。清洗这一步放在爬虫里做最简单因为数据规模还小pandas 就能搞定。下面的代码把 CSV 里的空价格行删除并按自然键去重import pandas as pd df pd.read_csv(products.csv) df df.dropna(subset[min_price, max_price, avg_price]) df df.drop_duplicates(subset[report_date, market, product_name, spec]) df df[(df[avg_price] 0) (df[avg_price] 1000)] df.to_csv(products_clean.csv, indexFalse, encodingutf-8-sig)这里最关键的是去重逻辑同一市场同一天对同一个品名可能有多条不同规格要把spec加进自然键否则会把「三级黄瓜」和「一级黄瓜」误删。均价过滤区间则用来处理明显异常值比如 0 元或上万元的错误记录。提示清洗后的数据单独存一份products_clean.csv原始 CSV 保留不动。这样 Spark 分析出问题时可以回到原始数据排查不用重新跑爬虫。3. 用Spark做农产品价格聚合分析并输出可复用结果3.1 为什么课程设计里要选 Spark 而不是继续用 pandas如果只分析几万条农产品数据pandas 完全够用而且写起来更顺手。但课程设计要体现「大数据处理」的能力Spark 的价值在于分析代码写好后本地local[*]模式能跑换到集群上同样能跑不需要改逻辑。用 Spark 做聚合还能自然引出分区数、shuffle、缓存调优这些答辩高频问题。另一个实用原因是结果可复用用 Spark 算出的统计结果写入中间文件JSON/CSVFlask 启动时只读中间文件绝不直接调 Spark。这个架构决定了你答辩时 Web 页面响应快、稳定性高。3.2 用 PySpark 计算价格排行与波动率核心代码下面这段代码读入清好的 CSV按农产品名称分组计算均价、最高价、最低价以及波动率。波动率定义为(最高价 - 最低价) / 平均价可以粗略反映该农产品在不同市场间的价格离散程度。from pyspark.sql import SparkSession from pyspark.sql.functions import mean, min, max, col, round, count, when spark SparkSession.builder \ .appName(AgriPriceAnalysis) \ .master(local[*]) \ .config(spark.sql.shuffle.partitions, 4) \ .getOrCreate() df spark.read.csv(products_clean.csv, headerTrue, inferSchemaTrue) df df.withColumn(avg_price, df[avg_price].cast(double)) \ .withColumn(min_price, df[min_price].cast(double)) \ .withColumn(max_price, df[max_price].cast(double)) stats df.groupBy(product_name).agg( count(*).alias(sample_count), mean(avg_price).alias(avg_price), min(min_price).alias(min_price), max(max_price).alias(max_price) ).withColumn( volatility, round((col(max_price) - col(min_price)) / col(avg_price), 4) ) stats stats.filter(col(sample_count) 5) stats.orderBy(col(avg_price).desc()).show(10) stats.write.option(header, True).csv(output/price_stats, modeoverwrite) spark.stop()spark.sql.shuffle.partitions设为 4 是为本地小数据量准备的如果在集群上跑默认 200 会建太多空任务增加调度开销。sample_count过滤是为了避免只有一两条报价的杂项品类跑进 Top 榜。写结果用csv(output/price_stats)而不是saveAsTable结果以目录形式产出后续 Flask 读取时按目录下的 part 文件处理即可。3.3 Spark 本地跑通的三个必调参数本地模式跑 Spark 不需要集群一条命令就能把上面的脚本跑起来spark-submit --master local[*] --driver-memory 2g analysis.py参数含义如下表参数推荐值说明--masterlocal[*]用本机全部 CPU 核数运行--driver-memory2g给 Driver 分配 2G 内存Python 和结果集都占这里--conf spark.sql.shuffle.partitions4~8本地数据量小分区太多会造成大量空任务常见报错是java.lang.OutOfMemoryError通常发生在--driver-memory给小了而数据集偏大时。课程设计数据量一般几十 MB2G 足够。更隐蔽的问题是 Windows 下 Spark 路径含中文报IllegalArgumentException时先检查SPARK_HOME和项目路径是否全英文。4. 用Flask把Spark结果封装成可视化Web接口4.1 Flask 里绝对不要现场调 Spark不少人会尝试在 Flask 路由里直接初始化 SparkSession然后run一圈数据分析再返回结果。这个做法有两个致命问题SparkSession 初始化要消耗数秒每次请求都初始化会让接口超时Spark 的 Executor 内存和 Flask 进程抢资源页面一打开整个应用都可能 OOM。提示正确架构是 Spark 离线跑批把结果写到磁盘Flask 启动时一次性读取进内存。这样 Spark 挂了不影响 Web 服务Web 服务重启也不用重算。4.2 写一个读取分析结果的 Flask 路由JSON 接口 模板渲染下面代码实现了两个接口/api/stats返回 Spark 算出的农产品均价 Top10/渲染一个带 ECharts 的页面。import json import glob import os from flask import Flask, jsonify, render_template app Flask(__name__) def load_stats(): result_files glob.glob(output/price_stats/part-*.csv) stats [] for path in result_files: with open(path, r, encodingutf-8) as f: lines [line.strip().split(,) for line in f.readlines()] for row in lines: stats.append({ product_name: row[0], sample_count: int(row[1]), avg_price: float(row[2]), min_price: float(row[3]), max_price: float(row[4]), volatility: float(row[5]) }) return stats STATS_CACHE load_stats() app.route(/api/stats) def api_stats(): return jsonify(STATS_CACHE) app.route(/) def index(): return render_template(index.html, statsSTATS_CACHE) if __name__ __main__: app.run(host0.0.0.0, port5000, debugFalse)load_stats()在模块加载时执行一次把 Spark 输出的多个 part 文件合并成 Python 对象后续接口直接读内存数组。比每次请求都重新读磁盘要快很多也避免多线程场景下反复打开文件句柄。debugFalse在演示时很重要Werkzeug 的调试器如果开着异常页面会泄露代码路径。4.3 可视化页面ECharts 展示 Spark 分析结果前端用 ECharts 是常见做法这种方式的好处是无需引入重量级框架一个 HTML 就能完成柱状图和趋势图。下面给出templates/index.html的核心片段!DOCTYPE html html head meta charsetutf-8 title农产品价格可视化分析/title script srchttps://cdn.jsdelivr.net/npm/echarts5.4.3/dist/echarts.min.js/script /head body div idchart stylewidth: 800px; height: 500px; margin: 40px auto;/div script fetch(/api/stats) .then(resp resp.json()) .then(data { const sorted data.sort((a, b) b.avg_price - a.avg_price).slice(0, 10); const chart echarts.init(document.getElementById(chart)); chart.setOption({ title: { text: 农产品均价 TOP10 }, xAxis: { type: category, data: sorted.map(item item.product_name) }, yAxis: { type: value, name: 元/斤 }, series: [{ type: bar, data: sorted.map(item item.avg_price) }] }); }); /script /body /html页面从/api/stats拉数据在前端排序取 TOP10。可以看到 Flask 层完全不关心数据怎么算的只关心结果长什么样。这样职责划分清晰答辩时你也能一句话说清Spark 算完落盘Flask 只负责读和展示。4.4 把系统跑通的完整命令串按以下顺序启动就能在浏览器localhost:5000看到可视化页面# 第一步爬取数据 python spider.py # 第二步清洗数据 python clean.py # 第三步Spark 分析 spark-submit --master local[*] --driver-memory 2g analysis.py # 第四步启动 Flask python app.py每步产物独立products_clean.csv是清洗结果output/price_stats是分析结果。因此你可以跳过爬虫直接用现成数据集跑后三步这在演示时间紧张时非常有用。5. 演示之前数据验证、参数调整与集群扩展方向5.1 验证 Spark 输出是否正确抽样比对原始数据答辩时最怕被问「你这个统计结果对不对」。建议演示前做一次抽样验证从 Spark 输出中挑三个农产品手工在原始 CSV 里算平均价看是否一致。具体方式可以用 Excel 筛选该品类的记录用内置公式计算均价再和 Spark 页面上的数字对比。这样比对比口头解释可靠得多。验证代码可以写成一个独立脚本用来做系统化校验import pandas as pd from pyspark.sql import SparkSession raw pd.read_csv(products_clean.csv) spark SparkSession.builder.master(local[*]).appName(verify).getOrCreate() result spark.read.csv(output/price_stats, headerTrue).toPandas() for product in result.head(3)[product_name].tolist(): expected raw[raw[product_name] product][avg_price].mean() actual result[result[product_name] product][avg_price].iloc[0] assert abs(expected - float(actual)) 0.01, f{product} 校验失败这里有个值得注意的差异点Spark 读取 CSV 时如果字段被推断为字符串计算顺序与 pandas 在不同版本上可能略有差异因此误差容忍度设为 0.01。5.2 Spark 内存调小技巧面对「内存不足」别急着加内存把spark.sql.shuffle.partitions调大或调小很多时候比加内存更有效。如果分组键的基数较低比如只有 30 种农产品那么分区设 4 个就够。但如果你加了「按市场分组」或者「按日期分组」聚合维度变多分区数可以调整到8~16。经验是把分区数控制在executor 核数 × 3以内。本地模式则只看spark.sql.shuffle.partitions这一项。另外大数据集跑 DataFrame API 时强烈建议在聚合前做一次数据裁剪df df.filter(col(report_date) 2024-01-01)这一步能显著降低 shuffle 数据量。答辩时你能说出这个优化点会让系统「大数据能力」的可信度提升不少。5.3 从本地到集群讲清楚配置变化课程设计如果能顺带展示 Spark 集群部署意识是把「高分项目」区分开来的加分项。常见的集群部署策略是在三台虚拟机或云主机上安装 Spark standalone 模式分析脚本不变只改--master参数spark-submit \ --master spark://node01:7077 \ --executor-memory 2g \ --total-executor-cores 6 \ analysis.py你需要解释的是本地local[*]和集群模式的区别在于 Driver 和 Executor 的物理分布集群模式下同一个SparkSession代码会自动把任务分发到各节点。这句话比背十个名词都管用。提示三节点集群如果用作演示建议至少给每个 Executor 分配 1G 内存并关闭 Web UI 之外的闲置端口避免 Demo 现场因为资源占用被卡死。5.4 答辩演示的高分细节把「数据流水线」讲清楚最后落在演示技巧上。进入系统页面后不要直接点图表先按下面三步走第一步打开爬虫采集的记录数截图说明数据来源可信第二步切到 Spark 作业日志或输出目录说明分析层确实执行了分布式任务第三步再展示 Flask 可视化页面讲指标业务含义。这种顺序让评委看到的是一个完整的数据工程链路而不是一个写死的页面。一个小技巧给 Flask 页面加一个时间刷新按钮后端读取output/price_stats目录的文件修改时间并显示在页脚。演示时提前重跑一次 Spark 作业页面上时间变化评委能直观感知「新数据进来了」。这个细节不用写复杂逻辑十行代码就能实现。本文还有配套的精品资源点击获取
上一篇/下一篇内容由系统自动关联
返回资讯列表 →