概述

本文基于《美妆订单项目与可视化》案例,详细讲解如何通过 Spark SQL 处理结构化订单数据,并结合 pyecharts 实现数据可视化。从项目搭建、数据清洗到高阶分析(如RFM模型),完整覆盖大数据分析的核心流程。以下是关键步骤与代码实现。


一、项目搭建与数据准备

  1. 环境配置

    • 使用 PyCharm 创建 BeautyProduct 项目,配置 Python 3.6 环境。

    • 确保 HDFS 服务启动,将美妆订单数据(beauty_prod_info.csv 和 beauty_prod_sales.csv)上传至 HDFS:

    bash

    hdfs dfs -mkdir -p /datas
    hdfs dfs -put beauty_prod*.csv /datas
  2. 数据加载与清洗

    • 通过 Spark SQL 读取 CSV 文件,定义 Schema 避免自动推断错误:

    python

    # 商品信息表
    df1 = spark.read.csv(
        "hdfs://localhost:9000/datas/beauty_prod_info.csv",
        schema='prod_id string, product string, cataB string, cataA string, price float'
    )
    # 订单数据表
    df2 = spark.read.csv(
        "hdfs://localhost:9000/datas/beauty_prod_sales.csv",
        schema="order_id string, od_date string, cust_id string, ..."
    )
    • 清洗数据:处理空值、重复行、非法字符(如“元”“个”),转换字段类型:

    python

    df_prod_sales = df2 \
        .withColumn("od_date", regexp_replace('od_date','#','-')) \
        .withColumn("od_quantity", regexp_replace('od_quantity','个','')) \
        .withColumn("od_price", col('od_price').cast('float')) \
        .filter("od_date < '2023-1-1'")

二、核心分析任务实现

1. 商品价格分析
  • 目标:找出每个商品小类中价格最高的前5种商品。

  • 实现:使用窗口函数 dense_rank() 按价格降序排名:

    python

    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
    """)
3. 地区与商品需求分析
  • 城市订购量 TOP20

    python

    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
    """)
  • 省份需求量统计

    python

    df_province_sale = spark.sql("""
        SELECT cust_province, SUM(od_quantity) AS total_quantity 
        FROM sales GROUP BY cust_province
    """)
4. 客户价值挖掘(RFM模型)
  • 步骤

    1. 统计客户最近购买时间、消费频率、总金额:

      python

      df = spark.sql("""
          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
      """)
    2. 归一化计算 R(Recency)、F(Frequency)、M(Monetary)得分:

      python

      df = spark.sql("""
          SELECT cust_id, 
              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 customer
      """)
    3. 综合得分:score = R*20% + F*30% + M*50%,筛选高价值客户。


三、数据可视化(pyecharts)

1. 月度销售趋势图
  • 代码

    python

    bar = Bar()
    bar.add_xaxis(月份列表)
    bar.add_yaxis("订购数量(万件)", y1)
    bar.add_yaxis("金额(亿元)", y2)
    bar.set_global_opts(title_opts=opts.TitleOpts(title="每月商品订购情况"))
    bar.render("month_sale.html")
  • 效果:双纵坐标柱状图展示销量与金额趋势。

2. 城市订购量 TOP20
  • 优化:横向柱状图提升可读性:

    python

    Bar()
      .add_yaxis("订购量", y, label_opts=opts.LabelOpts(position="right"))
      .reversal_axis()  # 翻转坐标轴
  • 效果:数据量大的城市显示在上方,直观对比需求差异。


四、结果保存与部署

  • 保存至 HDFS:将分析结果写入分布式存储,便于后续处理:

    python

    df_top5_product.write.csv("hdfs://localhost:9000/datas/result.top5_product")
    df_month_sale.write.csv("hdfs://localhost:9000/datas/result.month_sale")
  • 生产部署:通过 spark-submit 提交至 YARN 集群:

    bash

    spark-submit --master yarn --deploy-mode cluster main.py

五、总结

  • 技术栈:Spark SQL(数据处理)、窗口函数(复杂分析)、pyecharts(可视化)、HDFS(存储)。

  • 核心价值

    • 从原始数据到业务洞察的全流程实现。

    • RFM 模型助力精准客户分群。

    • 可视化图表直观呈现数据规律。


通过本案例,读者可掌握 Spark SQL 的核心操作 与 数据可视化 的实践方法,为实际业务中的大数据分析提供完整参考。

Logo

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

更多推荐