项目概述

本实战项目通过Spark SQL对美妆商品订单数据进行全面分析,并使用pyecharts实现数据可视化。项目完整展示了从数据清洗、预处理到多维分析、客户价值挖掘的完整流程,最终生成直观的图表展示分析结果。

环境准备与数据加载

1. 创建Spark项目

  • 使用PyCharm创建名为BeautyProduct的项目

  • 确保Python解释器版本为3.6

2. 数据准备

bash

# 在HDFS上创建数据目录并上传CSV文件
hdfs dfs -mkdir -p /datas
hdfs dfs -put beauty_prod*.csv /datas

3. 数据加载与Schema定义

python

# 商品信息表
df1 = spark.read.csv("hdfs://localhost:9000/datas/beauty_prod_info.csv",
    encoding='UTF-8', header=True, sep=',', inferSchema=False,
    schema='prod_id string, product string, cataB string, cataA string, price float')

# 订单销售表
df2 = spark.read.csv("hdfs://localhost:9000/datas/beauty_prod_sales.csv",
    encoding='UTF-8', header=True, sep=',', inferSchema=False,
    schema="""order_id string, od_date string, cust_id string,
    cust_region string, cust_province string, cust_city string,
    prod_id string, od_quantity string, od_price string, od_amount float""")

数据清洗与预处理

python

# 数据清洗流程
df_prod_info = df1.dropna().dropDuplicates(['prod_id'])

df_prod_sales = df2 \
    .dropna() \
    .distinct() \
    .withColumn("od_date", regexp_replace('od_date','#','-')) \
    .withColumn("od_quantity", regexp_replace('od_quantity','个','')) \
    .withColumn("od_price", regexp_replace('od_price','元','')) \
    .withColumn('cust_province', regexp_replace('cust_province','自治区|维吾尔|回族','')) \
    .withColumn("od_date", col('od_date').cast('date')) \
    .withColumn("od_quantity", col('od_quantity').cast('integer')) \
    .withColumn("od_price", col('od_price').cast('float')) \
    .withColumn("od_month", F.month('od_date')) \
    .filter("od_date<'2023-1-1'")

# 创建临时视图
df_prod_info.createOrReplaceTempView('prod')
df_prod_sales.createOrReplaceTempView('sales')

多维数据分析

1. 商品价格排名分析

python

# 各商品小类中价格排名前5的商品
df_top5_product = spark.sql("""
    SELECT *, dense_rank() over (partition by cataB order by price DESC) AS rank 
    FROM prod
""").where('rank<=5')

2. 月度销售趋势分析

python

# 每月商品订购数量和消费金额
df_month_sale = spark.sql("""
    SELECT od_month, sum(od_quantity) AS total_quantity,
    sum(od_amount) AS total_amount
    FROM sales
    GROUP BY od_month
    ORDER BY od_month
""")

3. 区域销售分析

python

# 订购量TOP20城市
df_city_sale = spark.sql("""
    SELECT cust_city, sum(od_quantity) AS total_quantity
    FROM sales
    GROUP BY cust_city
    ORDER BY total_quantity DESC
    LIMIT 20
""")

# 各省份订购量
df_province_sale = spark.sql("""
    SELECT cust_province, sum(od_quantity) AS total_quantity
    FROM sales
    GROUP BY cust_province
""")

4. 商品类别需求分析

python

# 各类商品订购量
df_best_seller = spark.sql("""
    SELECT p.cataA, p.cataB, sum(s.od_quantity) AS total_quantity
    FROM sales s JOIN prod p ON s.prod_id=p.prod_id
    GROUP BY p.cataA, p.cataB
""")

5. 客户价值分析(RFM模型)

python

# RFM模型客户价值分析
df_customerRFM = spark.sql("""
    SELECT cust_id, od_latest, total_count, total_amount,
    CASE WHEN R>0.8 AND F>0.8 AND M>0.8 THEN '高价值客户'
         WHEN R>0.5 AND F>0.5 AND M>0.5 THEN '潜力客户'
         ELSE '一般客户' END AS score
    FROM (
        SELECT cust_id, od_latest, total_count, total_amount,
        percent_rank() over (order by od_latest) as R,
        percent_rank() over (order by total_count) as F,
        percent_rank() over (order by total_amount) as M
        FROM (
            SELECT cust_id, max(od_date) AS od_latest,
            count(order_id) AS total_count,
            sum(od_amount) AS total_amount
            FROM sales GROUP BY cust_id
        )
    )
""")

分析结果保存

python

# 将各分析结果保存到HDFS
df_top5_product.write.csv('hdfs://localhost:9000/datas/result.top5_product')
df_month_sale.write.csv('hdfs://localhost:9000/datas/result.month_sale')
df_city_sale.write.csv('hdfs://localhost:9000/datas/result.city_sale')
df_best_seller.write.csv('hdfs://localhost:9000/datas/result.best_seller')
df_province_sale.write.csv('hdfs://localhost:9000/datas/result.province_sale')
df_customerRFM.write.csv('hdfs://localhost:9000/datas/result.customerRFM')

数据可视化实现

1. 月度销售趋势图

python

# 每月商品订购情况柱状图
c = (
    Bar()
    .add_xaxis(x)
    .add_yaxis("订购数量(万件)", y1)
    .add_yaxis("金额(亿元)", y2)
    .set_global_opts(
        title_opts=opts.TitleOpts(title="每月商品订购情况"),
        yaxis_opts=opts.AxisOpts(name="金额(亿元)"),
        xaxis_opts=opts.AxisOpts(name="月份")
    )
    .render("month_sale.html")
)

2. 城市销售TOP20横向柱状图

python

# 城市订购量排名TOP20
c = (
    Bar()
    .add_xaxis(x)
    .add_yaxis("订购量", y,
        label_opts=opts.LabelOpts(position="right",formatter='{@[1]}万'))
    .reversal_axis()
    .set_global_opts(
        title_opts=opts.TitleOpts("订购数量排名TOP20")
    )
    .render("city_sale.html")
)

项目总结

本实战项目完整展示了使用Spark SQL处理和分析电商订单数据的全流程:

  1. 数据准备:将CSV数据加载到HDFS,定义明确的数据Schema

  2. 数据清洗:处理缺失值、重复值、格式转换和异常值

  3. 多维分析:从商品、时间、地域、客户等多个维度深入分析

  4. 高级分析:应用RFM模型进行客户价值分层

  5. 可视化:使用pyecharts生成直观的交互式图表

通过本项目,读者可以掌握:

  • Spark SQL进行结构化数据处理的核心技巧

  • 电商数据分析的典型场景和方法

  • 使用窗口函数解决复杂分析问题

  • 基于RFM模型的客户价值分析方法

  • 数据分析结果的可视化呈现技巧

项目代码可直接应用于实际业务场景,也可根据具体需求进行扩展和优化。

Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐