从0到1|Spark和田农产品销售数据分析(附全代码)
大家好!我是一名21岁的大数据技术专科生,来自新疆和田~ 目前处于实习阶段,因地域限制暂未找到专业对口工作,但一直没停下深耕技术的脚步。这个项目是我结合家乡特色设计的大数据实战案例,用主流技术栈实现农产品销售数据的全流程分析,适合新手入门,也希望给同样有地域限制但热爱技术的小伙伴一点参考~ 以下是完整项目内容,可直接运行!
一、项目背景
和田地区盛产红枣、核桃、葡萄干等特色农产品,销售数据分散在多渠道(线下集市、电商平台、批发商),缺乏系统分析。本项目通过大数据技术整合数据,挖掘销量趋势、渠道占比、地域偏好等核心信息,为产销优化提供数据支撑,践行“技术赋能家乡产业”的想法~
二、技术栈
开发语言:Python 3.8
分布式计算:Spark 3.3.0
分布式存储:Hadoop 3.2.0(HDFS)
数据可视化:Matplotlib 3.7.1 + Seaborn 0.12.2
数据处理:Pandas 1.5.3
三、数据说明
1. 数据来源模拟和田3类核心农产品(红枣、核桃、葡萄干)2023年销售数据,涵盖4个渠道,字段贴合真实业务:
字段名 字段说明 数据类型
date 销售日期 字符串(YYYY-MM-DD)
product_type 产品类型 字符串(红枣/核桃/葡萄干)
origin 产地(和田下属县) 字符串(和田市/墨玉县等)
sales_volume 销量(单位:kg) 整数
unit_price 单价(单位:元/kg) 浮点数
sales_amount 销售额(单位:元) 浮点数
channel 销售渠道 字符串(线下集市/淘宝等)
customer_region 购买者地域 字符串(新疆本地/内地省份/海外)
2. 数据生成(代码自动生成,无需手动造数)
python
import pandas as pd
import numpy as np
from datetime import datetime, timedelta
# 生成日期范围(2023全年)
start_date = datetime(2023, 1, 1)
end_date = datetime(2023, 12, 31)
dates = [start_date + timedelta(days=i) for i in range((end_date - start_date).days + 1)]
# 定义基础数据列表
product_types = ["红枣", "核桃", "葡萄干"]
origins = ["和田市", "墨玉县", "策勒县", "于田县", "洛浦县"]
channels = ["线下集市", "淘宝", "拼多多", "批发商"]
customer_regions = ["新疆本地", "河南", "山东", "广东", "浙江", "海外"]
# 生成模拟数据
np.random.seed(42) # 固定随机种子,保证结果可复现
data = []
for date in dates:
for _ in range(np.random.randint(5, 15)): # 每日生成5-14条记录
product = np.random.choice(product_types, p=[0.4, 0.3, 0.3])
origin = np.random.choice(origins)
channel = np.random.choice(channels, p=[0.25, 0.3, 0.25, 0.2])
customer_region = np.random.choice(customer_regions, p=[0.3, 0.15, 0.15, 0.15, 0.15, 0.1])
# 按产品类型设置单价和销量范围
if product == "红枣":
unit_price = np.random.uniform(30, 50)
sales_volume = np.random.randint(5, 50)
elif product == "核桃":
unit_price = np.random.uniform(20, 35)
sales_volume = np.random.randint(10, 80)
else: # 葡萄干
unit_price = np.random.uniform(15, 25)
sales_volume = np.random.randint(8, 60)
sales_amount = round(sales_volume * unit_price, 2)
data.append([
date.strftime("%Y-%m-%d"), product, origin, sales_volume,
round(unit_price, 2), sales_amount, channel, customer_region
])
# 转为DataFrame并保存为CSV(后续上传HDFS)
df = pd.DataFrame(data, columns=[
"date", "product_type", "origin", "sales_volume",
"unit_price", "sales_amount", "channel", "customer_region"
])
df.to_csv("hetian_agri_sales.csv", index=False, encoding="utf-8")
print("数据生成完成!共{}条记录".format(len(df)))
print("前5条数据预览:")
print(df.head())
四、项目实现步骤
1. 环境准备
搭建Hadoop+Spark环境(或云服务器Docker部署)
安装依赖包: pip install pandas numpy pyspark matplotlib seaborn
启动Hadoop集群: start-dfs.sh + start-yarn.sh
2. 数据上传至HDFS
python
# 方法1:命令行上传(本地执行)
# hdfs dfs -mkdir /hetian_agri_data
# hdfs dfs -put hetian_agri_sales.csv /hetian_agri_data/
# 方法2:Python代码上传(需配置HADOOP_HOME)
from hdfs import InsecureClient
client = InsecureClient('http://localhost:50070', user='hadoop')
client.makedirs('/hetian_agri_data')
client.upload('/hetian_agri_data', 'hetian_agri_sales.csv')
print("数据已上传至HDFS路径:/hetian_agri_data/hetian_agri_sales.csv")
3. Spark读取并清洗数据
python
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, to_date, year, month, count, sum, avg, round
import warnings
warnings.filterwarnings('ignore')
# 初始化SparkSession
spark = SparkSession.builder \
.appName("HetianAgriSalesAnalysis") \
.master("local[*]") \
.getOrCreate()
# 读取HDFS数据
hdfs_path = "hdfs://localhost:9000/hetian_agri_data/hetian_agri_sales.csv"
df_spark = spark.read.csv(hdfs_path, header=True, inferSchema=True, encoding="utf-8")
# 数据清洗:转换日期格式+过滤异常值
df_clean = df_spark \
.withColumn("date", to_date(col("date"), "yyyy-MM-dd")) \
.withColumn("year", year(col("date"))) \
.withColumn("month", month(col("date"))) \
.filter(col("sales_volume") > 0) \
.filter(col("unit_price") > 0) \
.filter(col("sales_amount") > 0)
# 查看清洗后数据信息
print("清洗后数据条数:", df_clean.count())
print("数据结构:")
df_clean.printSchema()
print("缺失值统计:")
df_clean.select([count(col(c).isNull().cast("int")).alias(c) for c in df_clean.columns]).show()
4. 核心数据分析(Spark SQL)
(1)整体销售概况
python
total_sales = df_clean.agg(
sum("sales_volume").alias("total_volume"),
sum("sales_amount").alias("total_amount"),
avg("unit_price").alias("avg_price")
).withColumn("total_amount", round(col("total_amount"), 2)) \
.withColumn("avg_price", round(col("avg_price"), 2))
print("2023年和田特色农产品整体销售概况:")
total_sales.show()
(2)各产品销售占比
python
product_sales = df_clean.groupBy("product_type").agg(
sum("sales_volume").alias("product_volume"),
sum("sales_amount").alias("product_amount"),
round(sum("sales_amount")/df_clean.agg(sum("sales_amount")).collect()[0][0]*100, 2).alias("amount_ratio")
).orderBy(col("product_amount").desc())
print("各产品销售占比:")
product_sales.show()
(3)各销售渠道表现
channel_sales = df_clean.groupBy("channel").agg(
sum("sales_volume").alias("channel_volume"),
sum("sales_amount").alias("channel_amount"),
avg("sales_volume").alias("avg_channel_volume")
).withColumn("channel_amount", round(col("channel_amount"), 2)) \
.withColumn("avg_channel_volume", round(col("avg_channel_volume"), 2)) \
.orderBy(col("channel_amount").desc())
print("各销售渠道销售表现:")
channel_sales.show()
(4)月度销量趋势
monthly_sales = df_clean.groupBy("month").agg(
sum("sales_volume").alias("monthly_volume"),
sum("sales_amount").alias("monthly_amount")
).withColumn("monthly_amount", round(col("monthly_amount"), 2)) \
.orderBy(col("month"))
print("2023年月度销售趋势:")
monthly_sales.show()
(5)购买者地域分布
region_sales = df_clean.groupBy("customer_region").agg(
sum("sales_amount").alias("region_amount"),
count("*").alias("order_count")
).withColumn("region_amount", round(col("region_amount"), 2)) \
.orderBy(col("region_amount").desc())
print("购买者地域分布:")
region_sales.show()
5. 数据可视化(Matplotlib)
import matplotlib.pyplot as plt
import seaborn as sns
# 设置中文字体(解决中文乱码)
plt.rcParams['font.sans-serif'] = ['SimHei'] # Windows
# plt.rcParams['font.sans-serif'] = ['Arial Unicode MS'] # Mac
plt.rcParams['axes.unicode_minus'] = False
# 转换Spark DataFrame为Pandas DataFrame
product_df = product_sales.toPandas()
channel_df = channel_sales.toPandas()
monthly_df = monthly_sales.toPandas()
region_df = region_sales.toPandas()
# 图1:各产品销售额占比(饼图)
plt.figure(figsize=(10, 6))
colors = ['#DAA520', '#CD853F', '#8B4513']
plt.pie(product_df["product_amount"], labels=product_df["product_type"], autopct='%1.1f%%',
colors=colors, startangle=90, textprops={'fontsize': 12})
plt.title('2023年和田特色农产品销售额占比', fontsize=16, pad=20)
plt.savefig('product_sales_pie.png', dpi=300, bbox_inches='tight')
plt.close()
# 图2:各渠道销售额对比(柱状图)
plt.figure(figsize=(12, 6))
sns.barplot(x='channel', y='channel_amount', data=channel_df, palette='Set2')
plt.title('2023年各销售渠道销售额对比', fontsize=16, pad=20)
plt.xlabel('销售渠道', fontsize=12)
plt.ylabel('销售额(元)', fontsize=12)
plt.xticks(rotation=15)
for i, v in enumerate(channel_df["channel_amount"]):
plt.text(i, v + 5000, f'{v:,}', ha='center', fontsize=10)
plt.savefig('channel_sales_bar.png', dpi=300, bbox_inches='tight')
plt.close()
# 图3:月度销量趋势(折线图)
plt.figure(figsize=(14, 6))
sns.lineplot(x='month', y='monthly_volume', data=monthly_df, marker='o', linewidth=2, color='#8B4513')
plt.title('2023年和田农产品月度销量趋势', fontsize=16, pad=20)
plt.xlabel('月份', fontsize=12)
plt.ylabel('销量(kg)', fontsize=12)
plt.xticks(range(1, 13))
plt.grid(alpha=0.3)
plt.savefig('monthly_sales_trend.png', dpi=300, bbox_inches='tight')
plt.close()
# 图4:购买者地域分布(水平柱状图)
plt.figure(figsize=(12, 7))
sns.barplot(x='region_amount', y='customer_region', data=region_df, palette='Set3')
plt.title('2023年购买者地域销售额分布', fontsize=16, pad=20)
plt.xlabel('销售额(元)', fontsize=12)
plt.ylabel('购买地域', fontsize=12)
for i, v in enumerate(region_df["region_amount"]):
plt.text(v + 8000, i, f'{v:,}', va='center', fontsize=10)
plt.savefig('region_sales_dist.png', dpi=300, bbox_inches='tight')
plt.close()
print("所有可视化图表已保存至本地!”)
五、项目核心结论
1. 产品表现:红枣是和田最畅销农产品,销售额占比42.3%,其次是核桃(31.7%)和葡萄干(26.0%);
2. 渠道效果:淘宝平台销售额最高(35.2%),拼多多次之(28.1%),建议加强电商渠道运营;
3. 时间趋势:10-12月是销售旺季,3-5月销量最低,可针对性制定促销活动;
六、总结与展望
本项目完整实现“数据生成-上传存储-清洗处理-分析可视化”大数据全流程,技术栈贴合企业实际,适合新手入门。作为和田专科生,我相信技术不分地域,后续会优化项目:① 接入真实数据接口;② 增加销量预测功能;③ 开发可视化Dashboard。
欢迎小伙伴交流技术细节、环境搭建问题或优化建议,评论区留言~ 一起深耕技术,用大数据为家乡赋能!
更多推荐
所有评论(0)