1、项目介绍

技术栈
Python 语言、Spark 大数据计算框架、Hadoop 分布式存储、Hive 数据仓库、Django 后端框架、Vue 前端框架、selenium 爬虫技术、Echarts 可视化库、DataV 可视化工具、机器学习线性回归预测模型、淘宝数据源

功能模块

  • 商品销售数据分析大屏模块
  • 销量预测模块
  • Spark 大数据分析模块
  • 商品数据采集模块
  • 商品数据分析模块
  • 后台管理模块

项目介绍
本系统基于 Spark 与 Hadoop 大数据技术,结合 Django 与 Vue 框架,构建商品销售数据智能化采集、分析与预测平台。系统通过 selenium 爬虫从淘宝获取实时商品数据,经清洗后存储于 Hive 数据仓库。Spark 承担大规模数据处理与分析任务,前端借助 Echarts 与 DataV 实现可视化大屏,展示各类型销量占比、省市分布、店铺排行及商品词云图等核心指标。销量预测模块采用机器学习线性回归模型,用户选择商品种类、产地、价格等条件后输出预测结果。Spark 大数据分析模块提供代码编辑运行环境,支持商品数据的批量处理与预览,为电商运营提供数据驱动的决策支持。

2、项目界面

(1)商品销售数据分析大屏
该页面是商品销售数据可视化平台的首页大屏,展示商品核心销售数据指标,可呈现各类型销量占比、商品省市分布、优秀店铺和商品销量排行,还能查看销售与销售额占比,生成商品相关词云图,同时具备销量预测相关功能模块。
在这里插入图片描述

(2)机器学习线性回归预测模型销量预测
该页面是商品销售数据可视化平台的销量预测页,设有条件选择区域可选择商品种类、是否包邮、商品产地、价格等相关筛选条件,配备查看操作按钮,结果展示区域可直观呈现根据所选条件计算出的商品销量预测结果,操作流程简洁清晰。

在这里插入图片描述

(3)spark大数据分析
该页面是代码编辑运行界面,可编写基于Spark的商品数据处理代码,能创建Spark上下文对象,执行商品数据的展示操作,可查看代码运行后的进程状态与退出代码,还能看到数据相关字段及内容预览,同时提供类型提示工具的安装提示功能。
在这里插入图片描述

3、项目说明

一、技术栈简要说明
系统后端采用 Python 语言与 Django 框架构建,前端使用 Vue 框架实现交互界面。大数据处理基于 Spark 计算框架与 Hadoop 分布式存储,数据仓库选用 Hive 进行大规模数据管理。数据采集使用 selenium 爬虫技术抓取淘宝商品信息,可视化部分通过 Echarts 与 DataV 实现图表大屏展示,预测模块采用机器学习线性回归模型进行销量预测。

二、功能模块详细介绍
· 商品销售数据分析大屏模块
作为系统首页大屏,集中展示商品核心销售数据指标,包含各类型销量占比、商品省市分布、优秀店铺和商品销量排行、销售与销售额占比等关键数据,同时生成商品相关词云图,并提供销量预测功能模块入口,实现多维度数据的可视化总览。

· 销量预测模块
提供条件选择区域,用户可筛选商品种类、是否包邮、商品产地、价格等参数,点击查看按钮后系统基于线性回归模型输出对应商品的销量预测结果,操作流程简洁清晰,为库存管理与销售策略提供数据参考。

· Spark 大数据分析模块
提供代码编辑运行界面,支持编写基于 Spark 的商品数据处理代码。用户可创建 Spark 上下文对象,执行商品数据的展示与处理操作,查看代码运行后的进程状态与退出代码,预览数据相关字段及内容,同时提供类型提示工具的安装提示功能,满足大数据开发与调试需求。

· 商品数据采集模块
基于 selenium 爬虫技术,从淘宝电商平台实时抓取商品销售数据,涵盖商品种类、价格、销量、产地、包邮情况等关键字段,经清洗后存储至 Hive 数据仓库,为后续分析与预测提供稳定数据源。

· 商品数据分析模块
依托 Spark 与 Hadoop 对海量商品数据进行多维度分析,涵盖销量分布、地域特征、店铺表现、品类趋势等维度,挖掘数据内在规律,为可视化大屏与决策支持提供分析基础。

· 后台管理模块
提供系统配置与数据管理功能,支持用户权限管控、数据维护、任务调度等操作,保障平台稳定运行与数据准确性。

三、项目总结
本系统基于 Spark 与 Hadoop 大数据技术,结合 Django 与 Vue 框架,构建了商品销售数据智能化采集、分析与预测平台。系统通过 selenium 爬虫从淘宝获取实时商品数据,经清洗后存储于 Hive 数据仓库。Spark 承担大规模数据处理与分析任务,前端借助 Echarts 与 DataV 实现可视化大屏,展示各类型销量占比、省市分布、店铺排行及商品词云图等核心指标。销量预测模块采用机器学习线性回归模型,用户选择商品种类、产地、价格等条件后输出预测结果。Spark 大数据分析模块提供代码编辑运行环境,支持商品数据的批量处理与预览,为电商运营提供数据驱动的决策支持。

4、核心代码

# coding:utf8

#spark导包
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.functions import count
from pyspark.sql.functions import monotonically_increasing_id
from pyspark.sql.types import StructType,StructField,IntegerType,StringType,FloatType
from pyspark.sql.functions import col,sum,when
from pyspark.sql.functions import desc,asc

if __name__ == '__main__':
    #构建
    spark = SparkSession.builder.appName("sparkSQL").master("local[*]"). \
        config("spark.sql.shuffle.partitions", 2). \
        config("spark.sql.warehouse.dir", "hdfs://node1:8020/user/hive/warehouse"). \
        config("hive.metastore.uris", "thrift://node1:9083"). \
        enableHiveSupport(). \
        getOrCreate()

    sc = spark.sparkContext

    #读取Hive表
    commoditydata = spark.read.table('commoditydata')
    # commoditydata.show()

    #需求一 地址
    #以address字段为组,统计
    result1 = commoditydata.groupby("address").count()

    #将结果转化为DF对象
    df = result1.toPandas()
    # print(df)

    #sql
    result1.write.mode("overwrite").\
        format("jdbc").\
        option("url","jdbc:mysql://node1:3306/bigdata?useSSL=false&useUnicode=true&charset=utf8").\
        option("dbtable","addresssort").\
        option("user","root").\
        option("password","root").\
        option("encoding","utf-8").\
        save()
    #写入HIVE
    result1.write.mode("overwrite").saveAsTable("addresssort","parquet")
    # spark.sql("select * from addresssort").show()

    #需求二:各类型销量占比
    #以type为组
    result2 = commoditydata.groupBy("type").sum("buy_len").withColumnRenamed("sum(buy_len)","total_bue_len")

    #
    df = result2.toPandas()
    # print(df)

    #sql
    result2.write.mode("overwrite").\
        format("jdbc").\
        option("url","jdbc:mysql://node1:3306/bigdata?useSSL=false&useUnicode=true&charset=utf8").\
        option("dbtable","typesort").\
        option("user","root").\
        option("password","root").\
        option("encoding","utf-8").\
        save()
    #HIVE
    result2.write.mode("overwrite").saveAsTable("typesort","parquet")
    # spark.sql("select * from typesort").show()

    #需求三:统计每个省市日销100+的店铺 TOP10
    #province
    #按“address",”name“为组
    result3 = commoditydata.groupBy("address","name").\
        agg({"buy_len":"sum"}).withColumnRenamed("sum(buy_len)","total_buy_len").\
        filter("total_buy_len > 100").\
        drop_duplicates(subset=["name"]).\
        groupby("address").count()
    result3 = result3.orderBy(col("count").desc()).limit(10)
    df = result3.toPandas()
    # print(df)
    #sql
    result3.write.mode("overwrite").\
        format("jdbc").\
        option("url","jdbc:mysql://node1:3306/bigdata?useSSL=false&useUnicode=true&charset=utf8").\
        option("dbtable","provincesort").\
        option("user","root").\
        option("password","root").\
        option("encoding","utf-8").\
        save()
    #HIVE
    result3.write.mode("overwrite").saveAsTable("provincesort","parquet")
    # spark.sql("select * from provincesort").show()

    #需求四:统计每种类型buy_len的总和以及buy_len与price的总和
    result4 = commoditydata.groupBy("type").agg(
        F.sum("buy_len").alias("total_buy_len"),
        F.sum(F.col("buy_len") * F.col("price")).alias("sum_buy_len_price")
    )
    result4 = result4.withColumn("ration",F.col("sum_buy_len_price")/F.col("total_buy_len"))
    df = result4.toPandas()
    # print(df)
    #sql
    result4.write.mode("overwrite").\
        format("jdbc").\
        option("url","jdbc:mysql://node1:3306/bigdata?useSSL=false&useUnicode=true&charset=utf8").\
        option("dbtable","salerate").\
        option("user","root").\
        option("password","root").\
        option("encoding","utf-8").\
        save()
    #HIVE
    result4.write.mode("overwrite").saveAsTable("salerate","parquet")
    # spark.sql("select * from salerate").show()

    #需求五:店铺/商品销售额统计
    resultDouble = commoditydata.withColumn("volume",when(col("buy_len") * col("price") < 100,"小于100")
                                            .when(col("buy_len") * col("price") < 200,"小于200")
                                            .when(col("buy_len") * col("price") < 500,"小于500")
                                            .when(col("buy_len") * col("price") < 1000,"小于1000")
                                            .when(col("buy_len") * col("price") < 5000,"小于5000")
                                            .when(col("buy_len") * col("price") < 10000,"小于10000")
                                            .otherwise("大于10000"))
    result5 = resultDouble.drop_duplicates(["title"])
    result5 = result5.groupby("volume").agg(count("title").alias("title_count"))

    result6 = resultDouble.drop_duplicates(["name"])
    result6 = result6.groupby("volume").agg(count("name").alias("name_count"))

    # result5.show()
    # result6.show()
    result_cobine = result5.join(result6,"volume")
    # result_cobine.show()

    #sql
    result_cobine.write.mode("overwrite").\
        format("jdbc").\
        option("url","jdbc:mysql://node1:3306/bigdata?useSSL=false&useUnicode=true&charset=utf8").\
        option("dbtable","storeproduct").\
        option("user","root").\
        option("password","root").\
        option("encoding","utf-8").\
        save()
    #HIVE
    result_cobine.write.mode("overwrite").saveAsTable("storeproduct","parquet")
    # spark.sql("select * from storeproduct").show()

    #需求6:商品销量top10

    #productTop
    resultList = commoditydata.orderBy("buy_len",ascending=False).limit(10)
    # resultList.show()

    #sql
    resultList.write.mode("overwrite").\
        format("jdbc").\
        option("url","jdbc:mysql://node1:3306/bigdata?useSSL=false&useUnicode=true&charset=utf8").\
        option("dbtable","productTop10").\
        option("user","root").\
        option("password","root").\
        option("encoding","utf-8").\
        save()
    #HIVE
    resultList.write.mode("overwrite").saveAsTable("productTop10","parquet")
    spark.sql("select * from productTop10").show()



更多推荐