基于Hadoop MapReduce的阿迪达斯季度销量数据分析系统

一、项目概述

本项目是一个基于Hadoop MapReduce的分布式数据分析平台,专门用于处理和分析阿迪达斯的季度销售数据。系统采用大数据技术栈,实现了从数据采集、处理、存储到可视化展示的完整数据流,为业务决策提供多维度数据支持。

1.1 项目背景

随着全球运动品牌市场竞争日益激烈,阿迪达斯需要深入分析其季度销售数据以制定更有效的营销策略。传统的数据分析方法难以处理大规模历史数据,无法快速获取多维度分析结果。因此,我们构建了一个高效的分布式数据分析平台。

1.2 核心功能

  • 多维度数据分析:支持按年份、品牌、地区等维度进行统计分析
  • 品牌对比分析:实现阿迪达斯与耐克的销售数据对比
  • 经济指标关联:结合各地区GDP、价格指数等经济指标进行综合分析
  • 可视化展示:通过图表直观展示分析结果

二、技术架构

2.1 技术栈

层次技术选型版本
数据采集Flume-
数据存储Hadoop HDFS3.3.4
数据处理Hadoop MapReduce3.3.4
数据传输Sqoop-
数据存储MySQL8.0
后端服务Spring Boot2.7.18
ORM框架MyBatis Plus3.5.3.1
前端框架Vue 33.5.26
UI组件库Element Plus2.13.1
可视化库ECharts5.5.1
构建工具Vite7.3.0

2.2 项目结构

adsAnalysis/
├── ads-api/              # 后端RESTful API服务
│   ├── src/main/java/org/ads/api/
│   │   ├── controller/   # 控制器层
│   │   ├── service/      # 服务层
│   │   ├── mapper/       # 数据访问层
│   │   └── entity/       # 实体类
│   └── pom.xml
├── ads-mapreduce/        # Hadoop MapReduce数据处理模块
│   ├── src/main/java/org/ads/
│   │   ├── driver/       # 作业驱动类
│   │   ├── mapper/       # Mapper类
│   │   ├── reducer/      # Reducer类
│   │   └── util/         # 工具类
│   └── pom.xml
├── ads-dashboard/        # 前端可视化仪表盘
│   ├── src/
│   │   ├── views/        # 页面组件
│   │   ├── api/          # API接口
│   │   └── router/       # 路由配置
│   └── package.json
├── scripts/              # 脚本文件
│   ├── sqoop_export.sh   # Sqoop导出脚本
│   └── adsdata_created.sql  # 数据库建表脚本
└── pom.xml               # 父项目Maven配置

三、数据流向与处理流程

3.1 完整数据流

原始CSV数据 → Flume采集 → HDFS存储 → MapReduce处理 → MySQL存储 → Spring Boot API → Vue前端展示

3.2 详细处理步骤

  1. 数据采集:Flume实时采集阿迪达斯季度销量CSV数据
  2. 数据存储:数据存储到HDFS的/flume/data/01/ads_borrow目录
  3. 数据解析:使用AdidasSalesParser解析CSV数据
  4. MapReduce处理:运行6个MapReduce作业进行多维度分析
  5. 数据导出:使用Sqoop将分析结果导出到MySQL数据库
  6. API服务:Spring Boot提供RESTful API接口
  7. 前端展示:Vue 3 + ECharts实现数据可视化

四、核心代码解析

4.1 数据解析器

AdidasSalesParser负责解析CSV格式的原始数据:

public class AdidasSalesParser {
    private String date;           // 季度(如2000Q1)
    private double revenue;        // 阿迪达斯营收
    private double usGdp;          // 美国GDP
    private double europeGdp;      // 欧洲GDP
    private double chnGdpUs;       // 中国GDP(美元)
    private double priceIndex;     // 价格指数
    private double nike;           // 耐克营收
    private int year;              // 年份
    private int quarter;           // 季度

    public AdidasSalesParser(String line) {
        String[] fields = line.split(",");
        if (fields.length >= 7) {
            try {
                this.date = fields[0].trim();
                this.revenue = Double.parseDouble(fields[1].trim());
                this.usGdp = Double.parseDouble(fields[2].trim());
                this.europeGdp = Double.parseDouble(fields[3].trim());
                this.chnGdpUs = Double.parseDouble(fields[4].trim());
                this.priceIndex = Double.parseDouble(fields[5].trim());
                this.nike = Double.parseDouble(fields[6].trim());
                
                // 解析年份和季度
                if (this.date.matches("\\d{4}Q\\d")) {
                    this.year = Integer.parseInt(this.date.substring(0, 4));
                    this.quarter = Integer.parseInt(this.date.substring(5, 6));
                }
            } catch (NumberFormatException e) {
                // 异常处理:设置默认值
                this.revenue = 0.0;
                this.year = 0;
                this.quarter = 0;
            }
        }
    }
}

4.2 品牌总营收Mapper

统计Adidas和Nike的总营收:

public class BrandTotalRevenueMapper extends Mapper<Object, Text, Text, DoubleWritable> {
    private Text brand = new Text();
    private DoubleWritable revenue = new DoubleWritable();

    @Override
    protected void map(Object key, Text value, Context context) 
            throws IOException, InterruptedException {
        AdidasSalesParser parser = new AdidasSalesParser(value.toString());
        
        // 输出阿迪达斯营收
        brand.set("Adidas");
        revenue.set(parser.getRevenue());
        context.write(brand, revenue);
        
        // 输出耐克营收
        brand.set("Nike");
        revenue.set(parser.getNike());
        context.write(brand, revenue);
    }
}

4.3 品牌总营收Reducer

聚合计算各品牌总营收:

public class BrandTotalRevenueReducer extends Reducer<Text, DoubleWritable, Text, DoubleWritable> {
    private DoubleWritable totalRevenue = new DoubleWritable();

    @Override
    protected void reduce(Text key, Iterable<DoubleWritable> values, 
            Context context) throws IOException, InterruptedException {
        double sum = 0;
        for (DoubleWritable value : values) {
            sum += value.get();
        }

        // 格式化总营收,保留两位小数
        DecimalFormat df = new DecimalFormat("0.00");
        String formattedTotalRevenue = df.format(sum);
        double roundedTotalRevenue = Double.parseDouble(formattedTotalRevenue);

        totalRevenue.set(roundedTotalRevenue);
        context.write(key, totalRevenue);
    }
}

4.4 按年份和品牌统计Mapper

实现年份与品牌的复合维度分析:

public class YearBrandRevenueMapper extends Mapper<LongWritable, Text, Text, DoubleWritable> {
    private Text yearBrand = new Text();
    private DoubleWritable revenue = new DoubleWritable();

    @Override
    protected void map(LongWritable key, Text value, Context context) 
            throws IOException, InterruptedException {
        // 跳过表头行
        if (key.get() == 0) {
            return;
        }

        AdidasSalesParser parser = new AdidasSalesParser(value.toString());
        int year = parser.getYear();
        
        if (year > 0) {
            // 输出阿迪达斯营收
            String compositeKey = String.valueOf(year) + "\tAdidas";
            yearBrand.set(compositeKey);
            revenue.set(parser.getRevenue());
            context.write(yearBrand, revenue);
            
            // 输出耐克营收
            compositeKey = String.valueOf(year) + "\tNike";
            yearBrand.set(compositeKey);
            revenue.set(parser.getNike());
            context.write(yearBrand, revenue);
        }
    }
}

4.5 MapReduce作业驱动类

AdidasAnalysisDriver统一管理6个MapReduce作业:

public class AdidasAnalysisDriver {
    public static void main(String[] args) throws Exception {
        String inputPath = args.length > 0 ? args[0] : 
            "hdfs://192.168.199.101:8020/flume/data/01/ads_borrow";
        String baseOutputDir = "hdfs://192.168.199.101:8020/flume/output/";

        // 依次运行6个MapReduce作业
        boolean success = true;
        success &= runTeaTypeAvgPriceJob(inputPath, baseOutputDir + "avg_revenue_output");
        success &= runOriginTotalProductionJob(inputPath, baseOutputDir + "region_gdp_output");
        success &= runTeaTypeTotalSalesJob(inputPath, baseOutputDir + "brand_total_revenue_output");
        success &= runYearOriginTotalProductionJob(inputPath, baseOutputDir + "year_region_gdp_output");
        success &= runYearTeaTypeAvgPriceJob(inputPath, baseOutputDir + "year_avg_price_index_output");
        success &= runYearTeaTypeTotalSalesJob(inputPath, baseOutputDir + "year_brand_total_revenue_output");

        if (success) {
            System.out.println("所有作业执行成功!");
        } else {
            System.exit(-1);
        }
    }

    private static boolean runTeaTypeTotalSalesJob(String inputPath, String outputPath) 
            throws Exception {
        Configuration conf = createConfiguration();
        
        Path outputDir = new Path(outputPath);
        FileSystem fs = FileSystem.get(conf);
        if (fs.exists(outputDir)) {
            fs.delete(outputDir, true);
        }
        
        Job job = Job.getInstance(conf, "Brand Total Revenue");
        job.setJarByClass(AdidasAnalysisDriver.class);
        job.setMapperClass(BrandTotalRevenueMapper.class);
        job.setReducerClass(BrandTotalRevenueReducer.class);
        job.setOutputKeyClass(Text.class);
        job.setOutputValueClass(DoubleWritable.class);
        
        FileInputFormat.addInputPath(job, new Path(inputPath));
        FileOutputFormat.setOutputPath(job, outputDir);
        
        return job.waitForCompletion(true);
    }
}

4.6 Spring Boot API控制器

提供RESTful API接口:

@RestController
@RequestMapping("/api/brand-total-revenue")
public class BrandTotalRevenueController {
    @Autowired
    private BrandTotalRevenueService brandTotalRevenueService;

    @GetMapping("/all")
    public List<BrandTotalRevenue> getAllBrandTotalRevenue() {
        return brandTotalRevenueService.list();
    }
}

4.7 Vue前端可视化组件

使用ECharts展示品牌总营收:

<template>
  <div class="brand-total-revenue-container">
    <h2>品牌总营收</h2>
    <el-card class="chart-card">
      <div class="chart-container" ref="chartRef"></div>
    </el-card>
  </div>
</template>

<script setup>
import { ref, onMounted } from 'vue';
import * as echarts from 'echarts';
import { adsApi } from '../api';

const chartRef = ref(null);
let chartInstance = null;

const fetchDataAndRenderChart = async () => {
  const data = await adsApi.getBrandTotalRevenue();
  
  if (chartRef.value) {
    chartInstance = echarts.init(chartRef.value);
    const option = {
      title: { text: '品牌总营收排名', left: 'center' },
      tooltip: { trigger: 'axis', formatter: '{b}: {c} 亿元' },
      xAxis: { type: 'value', name: '总营收(亿元)' },
      yAxis: { 
        type: 'category', 
        data: data.map(item => item.brand).reverse() 
      },
      series: [{
        name: '总营收',
        type: 'bar',
        data: data.map(item => item.totalRevenue).reverse(),
        itemStyle: { color: '#fa8c16' }
      }]
    };
    chartInstance.setOption(option);
  }
};

onMounted(() => {
  fetchDataAndRenderChart();
});
</script>

五、数据库设计

5.1 数据库表结构

-- 按年份统计平均营收表
CREATE TABLE year_avg_revenue (
    id INT AUTO_INCREMENT PRIMARY KEY,
    year INT NOT NULL,
    avg_revenue DOUBLE NOT NULL,
    UNIQUE KEY uk_year (year)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

-- 各地区GDP总和表
CREATE TABLE region_gdp (
    id INT AUTO_INCREMENT PRIMARY KEY,
    region VARCHAR(50) NOT NULL,
    total_gdp DOUBLE NOT NULL,
    UNIQUE KEY uk_region (region)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

-- 品牌总营收表
CREATE TABLE brand_total_revenue (
    id INT AUTO_INCREMENT PRIMARY KEY,
    brand VARCHAR(50) NOT NULL,
    total_revenue DOUBLE NOT NULL,
    UNIQUE KEY uk_brand (brand)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

-- 按年份统计各地区GDP总和表
CREATE TABLE year_region_gdp (
    id INT AUTO_INCREMENT PRIMARY KEY,
    year INT NOT NULL,
    region VARCHAR(50) NOT NULL,
    total_gdp DOUBLE NOT NULL,
    UNIQUE KEY uk_year_region (year, region)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

-- 按年份统计平均价格指数表
CREATE TABLE year_avg_price_index (
    id INT AUTO_INCREMENT PRIMARY KEY,
    year INT NOT NULL,
    avg_price_index DOUBLE NOT NULL,
    UNIQUE KEY uk_year (year)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

-- 按年份统计各品牌总营收表
CREATE TABLE year_brand_total_revenue (
    id INT AUTO_INCREMENT PRIMARY KEY,
    year INT NOT NULL,
    brand VARCHAR(50) NOT NULL,
    total_revenue DOUBLE NOT NULL,
    UNIQUE KEY uk_year_brand (year, brand)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

六、Sqoop数据导出脚本

6.1 sqoop_export.sh脚本

#!/bin/bash

MYSQL_HOST="192.168.199.101"
MYSQL_PORT="3306"
MYSQL_DB="adsdata"
MYSQL_USER="root"
MYSQL_PASS="123456"
HDFS_BASE_DIR="/flume/output"

# 导出品牌总营收数据
sqoop export \
  --connect "jdbc:mysql://${MYSQL_HOST}:${MYSQL_PORT}/${MYSQL_DB}?useUnicode=true&characterEncoding=utf-8" \
  --username ${MYSQL_USER} \
  --password ${MYSQL_PASS} \
  --table brand_total_revenue \
  --columns "brand,total_revenue" \
  --export-dir ${HDFS_BASE_DIR}/brand_total_revenue_output \
  --input-fields-terminated-by "\t" \
  --input-lines-terminated-by '\n' \
  --num-mappers 1 \
  --batch \
  --mapreduce-job-name "BrandTotalRevenueExport"

# 导出其他表的数据...

七、项目部署与运行

7.1 环境要求

  • JDK:1.8+
  • Hadoop:3.3.4+
  • MySQL:8.0+
  • Node.js:20.19.0+
  • Maven:3.6+

7.2 构建MapReduce模块

cd ads-mapreduce
mvn clean package

7.3 运行MapReduce作业

hadoop jar ads-mapreduce-1.0-SNAPSHOT.jar org.ads.driver.AdidasAnalysisDriver

7.4 执行Sqoop导出

chmod +x scripts/sqoop_export.sh
./scripts/sqoop_export.sh

7.5 启动API服务

cd ads-api
mvn spring-boot:run

7.6 启动前端服务

cd ads-dashboard
npm install
npm run dev

八、数据分析维度

8.1 六大分析维度

分析维度Mapper类Reducer类说明
按年份统计平均营收YearAvgRevenueMapperYearAvgRevenueReducer计算每年平均营收
各地区GDP总和RegionGdpMapperRegionGdpReducer统计美国、欧洲、中国GDP
品牌总营收BrandTotalRevenueMapperBrandTotalRevenueReducer对比Adidas和Nike总营收
按年份统计各地区GDPYearRegionGdpMapperYearRegionGdpReducer每年各地区GDP数据
按年份统计平均价格指数YearAvgPriceIndexMapperYearAvgPriceIndexReducer每年平均价格指数
按年份统计各品牌总营收YearBrandRevenueMapperYearBrandRevenueReducer每年各品牌总营收

8.2 数据格式

原始数据格式

date,revenue,usGdp,europeGdp,chnGdpUs,priceIndex,nike
2023Q1,123456.78,25000.00,18000.00,12000.00,105.5,98765.43

MapReduce输出格式

Adidas	393703.57
Nike	297899.22

九、项目特色

9.1 技术亮点

  1. 分布式处理:基于Hadoop MapReduce实现大规模数据的分布式处理
  2. 模块化设计:每个分析维度独立实现,便于扩展和维护
  3. 完整数据流:从数据采集到可视化展示的完整数据处理链路
  4. 多维度分析:支持年份、品牌、地区等多种维度的统计分析
  5. 实时可视化:Vue 3 + ECharts实现动态数据可视化

9.2 性能优化

  1. 数据本地化:利用Hadoop的数据本地化特性减少网络传输
  2. 异常处理:完善的数据预处理机制确保系统稳定性
  3. 批量导出:Sqoop批量导出提高数据传输效率
  4. 缓存机制:前端数据缓存减少API调用次数

十、总结

本项目成功构建了一个基于Hadoop MapReduce的阿迪达斯季度销量数据分析系统,实现了从数据采集、处理、存储到可视化展示的完整数据流。系统采用分层架构设计,具有良好的可扩展性和维护性。

通过多维度数据分析,为业务决策提供了有力的数据支持。前端可视化界面直观展示了分析结果,帮助用户快速理解数据背后的业务价值。

未来可以进一步扩展系统的功能,如增加实时数据分析、引入机器学习算法进行销售预测、支持更多品牌的数据对比等。


作者:[大数据基础]
发布时间:2026年2月

相关文章

更多推荐