基于Hadoop MapReduce的阿迪达斯季度销量数据分析系统,大数据分析系统,flume上传HDFS,mapreduce分析,源码
·
基于Hadoop MapReduce的阿迪达斯季度销量数据分析系统
一、项目概述
本项目是一个基于Hadoop MapReduce的分布式数据分析平台,专门用于处理和分析阿迪达斯的季度销售数据。系统采用大数据技术栈,实现了从数据采集、处理、存储到可视化展示的完整数据流,为业务决策提供多维度数据支持。
1.1 项目背景
随着全球运动品牌市场竞争日益激烈,阿迪达斯需要深入分析其季度销售数据以制定更有效的营销策略。传统的数据分析方法难以处理大规模历史数据,无法快速获取多维度分析结果。因此,我们构建了一个高效的分布式数据分析平台。
1.2 核心功能
- 多维度数据分析:支持按年份、品牌、地区等维度进行统计分析
- 品牌对比分析:实现阿迪达斯与耐克的销售数据对比
- 经济指标关联:结合各地区GDP、价格指数等经济指标进行综合分析
- 可视化展示:通过图表直观展示分析结果
二、技术架构
2.1 技术栈
| 层次 | 技术选型 | 版本 |
|---|---|---|
| 数据采集 | Flume | - |
| 数据存储 | Hadoop HDFS | 3.3.4 |
| 数据处理 | Hadoop MapReduce | 3.3.4 |
| 数据传输 | Sqoop | - |
| 数据存储 | MySQL | 8.0 |
| 后端服务 | Spring Boot | 2.7.18 |
| ORM框架 | MyBatis Plus | 3.5.3.1 |
| 前端框架 | Vue 3 | 3.5.26 |
| UI组件库 | Element Plus | 2.13.1 |
| 可视化库 | ECharts | 5.5.1 |
| 构建工具 | Vite | 7.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 详细处理步骤
- 数据采集:Flume实时采集阿迪达斯季度销量CSV数据
- 数据存储:数据存储到HDFS的
/flume/data/01/ads_borrow目录 - 数据解析:使用AdidasSalesParser解析CSV数据
- MapReduce处理:运行6个MapReduce作业进行多维度分析
- 数据导出:使用Sqoop将分析结果导出到MySQL数据库
- API服务:Spring Boot提供RESTful API接口
- 前端展示: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类 | 说明 |
|---|---|---|---|
| 按年份统计平均营收 | YearAvgRevenueMapper | YearAvgRevenueReducer | 计算每年平均营收 |
| 各地区GDP总和 | RegionGdpMapper | RegionGdpReducer | 统计美国、欧洲、中国GDP |
| 品牌总营收 | BrandTotalRevenueMapper | BrandTotalRevenueReducer | 对比Adidas和Nike总营收 |
| 按年份统计各地区GDP | YearRegionGdpMapper | YearRegionGdpReducer | 每年各地区GDP数据 |
| 按年份统计平均价格指数 | YearAvgPriceIndexMapper | YearAvgPriceIndexReducer | 每年平均价格指数 |
| 按年份统计各品牌总营收 | YearBrandRevenueMapper | YearBrandRevenueReducer | 每年各品牌总营收 |
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 技术亮点
- 分布式处理:基于Hadoop MapReduce实现大规模数据的分布式处理
- 模块化设计:每个分析维度独立实现,便于扩展和维护
- 完整数据流:从数据采集到可视化展示的完整数据处理链路
- 多维度分析:支持年份、品牌、地区等多种维度的统计分析
- 实时可视化:Vue 3 + ECharts实现动态数据可视化
9.2 性能优化
- 数据本地化:利用Hadoop的数据本地化特性减少网络传输
- 异常处理:完善的数据预处理机制确保系统稳定性
- 批量导出:Sqoop批量导出提高数据传输效率
- 缓存机制:前端数据缓存减少API调用次数
十、总结
本项目成功构建了一个基于Hadoop MapReduce的阿迪达斯季度销量数据分析系统,实现了从数据采集、处理、存储到可视化展示的完整数据流。系统采用分层架构设计,具有良好的可扩展性和维护性。
通过多维度数据分析,为业务决策提供了有力的数据支持。前端可视化界面直观展示了分析结果,帮助用户快速理解数据背后的业务价值。
未来可以进一步扩展系统的功能,如增加实时数据分析、引入机器学习算法进行销售预测、支持更多品牌的数据对比等。
作者:[大数据基础]
发布时间:2026年2月
相关文章:
更多推荐
所有评论(0)