【大数据实战】基于Hadoop+Hive+SpringBoot+Vue的肥胖风险数据分析系统,hive数据etl分析,sqoop导数mysql
·
基于Hadoop+Hive+SpringBoot+Vue的肥胖风险数据分析系统
本文介绍了一个基于大数据技术栈的数据分析系统,实现了从数据采集、清洗、存储到可视化展示的完整流程。
一、项目概述
1.1 项目背景
肥胖已成为全球性的健康问题。本项目基于肥胖风险数据集,构建了一个数据分析系统,通过多维度数据分析,为健康管理和政策制定提供数据支持。
1.2 技术栈
大数据处理层:
- Hadoop 3.x:分布式存储和计算框架
- MapReduce:分布式数据处理
- Hive:数据仓库和SQL查询引擎
- Sqoop:数据迁移工具
后端服务层:
- Spring Boot 2.x:应用框架
- MyBatis-Plus:ORM框架
- MySQL:关系型数据库
前端展示层:
- Vue.js 2.x:前端框架
- ECharts:数据可视化库
- Element UI:UI组件库
1.3 项目架构
┌─────────────────────────────────────────────────────────┐
│ 前端展示层 (Vue.js) │
│ Dashboard | 趋势分析 | 风险因素 | 人口统计 | 生活方式 │
└─────────────────────────────────────────────────────────┘
↓ HTTP API
┌─────────────────────────────────────────────────────────┐
│ 后端服务层 (Spring Boot) │
│ ObesityAnalysisController + Service │
└─────────────────────────────────────────────────────────┘
↓ JDBC
┌─────────────────────────────────────────────────────────┐
│ 数据存储层 (MySQL) │
│ ADS应用层数据表 (ads_obesity_trend_analysis等) │
└─────────────────────────────────────────────────────────┘
↑ Sqoop
┌─────────────────────────────────────────────────────────┐
│ 数据仓库层 (Hive) │
│ ODS原始层 | DWD明细层 | DWS汇总层 | ADS应用层 │
└─────────────────────────────────────────────────────────┘
↑ MapReduce
┌─────────────────────────────────────────────────────────┐
│ 数据存储层 (HDFS) │
│ /obesity_output (清洗后的数据) │
└─────────────────────────────────────────────────────────┘
二、数据仓库设计
2.1 分层架构
项目采用数据仓库分层架构:
ODS层(原始数据层):
ods_obesity_raw:存储从MapReduce清洗后的原始数据
DWD层(明细数据层):
dwd_gender_info:性别维度表dwd_age_group_info:年龄分组维度表dwd_obesity_level_info:肥胖水平维度表dwd_family_history_info:家族史维度表dwd_lifestyle_info:生活方式维度表dwd_transport_info:交通方式维度表dwd_obesity_risk_fact:肥胖风险事实表
DWS层(汇总数据层):
dws_gender_obesity:性别肥胖分析表dws_age_group_obesity:年龄组肥胖分析表dws_obesity_level_distribution:肥胖水平分布表dws_lifestyle_obesity:生活方式肥胖分析表dws_family_history_obesity:家族史肥胖分析表dws_transport_obesity:交通方式肥胖分析表dws_bmi_distribution:BMI分布分析表
ADS层(应用数据层):
ads_obesity_trend_analysis:肥胖趋势分析表ads_obesity_risk_factors:风险因素分析表ads_obesity_demographic_analysis:人口统计学分析表ads_obesity_lifestyle_analysis:生活方式分析表ads_obesity_comprehensive_analysis:综合分析表
2.2 核心表结构
肥胖趋势分析表(ads_obesity_trend_analysis):
CREATE TABLE IF NOT EXISTS ads_obesity_trend_analysis (
age_group STRING COMMENT '年龄分组',
gender STRING COMMENT '性别',
total_count BIGINT COMMENT '总人数',
obesity_count BIGINT COMMENT '肥胖人数',
obesity_rate DOUBLE COMMENT '肥胖率(%)',
avg_bmi DOUBLE COMMENT '平均BMI',
avg_risk_score DOUBLE COMMENT '平均风险评分'
)
COMMENT '肥胖趋势分析表'
PARTITIONED BY (dt STRING COMMENT '日期分区')
ROW FORMAT DELIMITED
FIELDS TERMINATED BY '\001'
STORED AS TEXTFILE;
三、核心功能实现
3.1 MapReduce数据清洗
3.1.1 Driver类设计
ObesityDataDriver.java 是MapReduce作业的驱动类,负责作业的配置和提交。
核心代码:
public class ObesityDataDriver {
public static void main(String[] args) throws Exception {
String inputPath = "data/obesity_level.csv";
String outputPath = "/obesity_output";
Configuration conf = new Configuration();
conf.set("HADOOP_USER_NAME", "root");
conf.set("fs.hdfs.impl", "org.apache.hadoop.hdfs.DistributedFileSystem");
conf.set("fs.file.impl", "org.apache.hadoop.fs.LocalFileSystem");
UserGroupInformation.setConfiguration(conf);
UserGroupInformation ugi = UserGroupInformation.createRemoteUser("root");
UserGroupInformation.setLoginUser(ugi);
boolean useHDFS = false;
try {
conf.set("fs.defaultFS", "hdfs://192.168.199.101:8020");
FileSystem fs = FileSystem.get(conf);
fs.close();
useHDFS = true;
System.out.println("HDFS连接成功,将使用HDFS输出");
} catch (Exception e) {
System.out.println("HDFS连接失败: " + e.getMessage());
System.out.println("将使用本地文件系统");
useHDFS = false;
}
Job job = Job.getInstance(conf, "Obesity Data Cleaner");
job.setJarByClass(ObesityDataDriver.class);
job.setMapperClass(ObesityDataCleaner.ObesityCleanerMapper.class);
job.setReducerClass(ObesityDataCleaner.ObesityCleanerReducer.class);
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(NullWritable.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(NullWritable.class);
boolean success = job.waitForCompletion(true);
System.exit(success ? 0 : 1);
}
}
实现说明:
- HDFS连接检测:检测HDFS是否可用,失败时使用本地文件系统
- 用户权限管理:使用UserGroupInformation进行HDFS用户认证
- 路径配置:支持命令行参数动态配置输入输出路径
3.1.2 Mapper类设计
ObesityDataCleaner.ObesityCleanerMapper 负责数据的清洗和转换。
核心代码:
public static class ObesityCleanerMapper extends Mapper<Object, Text, Text, NullWritable> {
@Override
protected void map(Object key, Text value, Context context)
throws IOException, InterruptedException {
String line = value.toString().trim();
if (line.isEmpty()) {
return;
}
String[] fields = line.split(",");
if (fields.length < 17) {
context.getCounter("DataQuality", "InvalidFieldCount").increment(1);
return;
}
try {
String cleanedLine = cleanAndValidateRecord(fields, context);
if (cleanedLine != null) {
context.write(new Text(cleanedLine), NullWritable.get());
context.getCounter("DataQuality", "ValidRecords").increment(1);
} else {
context.getCounter("DataQuality", "FilteredRecords").increment(1);
}
} catch (Exception e) {
context.getCounter("DataQuality", "ProcessingErrors").increment(1);
}
}
}
数据清洗逻辑:
private String cleanAndValidateRecord(String[] fields, Context context) {
try {
String genderVal = normalizeGender(fields[gender]);
double ageVal = Double.parseDouble(fields[age].trim());
double heightVal = Double.parseDouble(fields[height].trim());
double weightVal = Double.parseDouble(fields[weight].trim());
if (heightVal < 1.0 || heightVal > 2.5) {
context.getCounter("DataQuality", "InvalidHeight").increment(1);
return null;
}
if (weightVal < 30 || weightVal > 200) {
context.getCounter("DataQuality", "InvalidWeight").increment(1);
return null;
}
if (ageVal < 14 || ageVal > 80) {
context.getCounter("DataQuality", "InvalidAge").increment(1);
return null;
}
double bmi = calculateBMI(weightVal, heightVal);
String ageGroup = calculateAgeGroup(ageVal);
int obesityLevelCode = getObesityLevelCode(obesityLevelVal);
double riskScore = calculateRiskScore(favcVal, fcvcVal, fafVal, smokeVal, familyHistoryVal);
StringBuilder sb = new StringBuilder();
sb.append(idVal).append(",");
sb.append(genderVal).append(",");
sb.append(df.format(ageVal)).append(",");
sb.append(df.format(heightVal)).append(",");
sb.append(df.format(weightVal)).append(",");
sb.append(familyHistoryVal).append(",");
sb.append(favcVal).append(",");
sb.append(df.format(fcvcVal)).append(",");
sb.append(df.format(ncpVal)).append(",");
sb.append(caecVal).append(",");
sb.append(smokeVal).append(",");
sb.append(df.format(ch2oVal)).append(",");
sb.append(sccVal).append(",");
sb.append(df.format(fafVal)).append(",");
sb.append(df.format(tueVal)).append(",");
sb.append(calcVal).append(",");
sb.append(mtransVal).append(",");
sb.append(obesityLevelVal).append(",");
sb.append(df.format(bmi)).append(",");
sb.append(ageGroup).append(",");
sb.append(obesityLevelCode).append(",");
sb.append(df.format(riskScore));
return sb.toString();
} catch (NumberFormatException e) {
context.getCounter("DataQuality", "ParseError").increment(1);
return null;
}
}
BMI计算和风险评分:
private double calculateBMI(double weight, double height) {
return weight / (height * height);
}
private double calculateRiskScore(int favc, double fcvc, double faf, int smoke, int familyHistory) {
double score = 0;
score += (favc == 1) ? 20 : 0;
score += Math.max(0, (3 - fcvc)) * 10;
score += Math.max(0, (3 - faf)) * 15;
score += (smoke == 1) ? 10 : 0;
score += (familyHistory == 1) ? 25 : 0;
return Math.min(100, Math.max(0, score));
}
实现说明:
- 数据验证:对身高、体重、年龄等字段进行校验
- 数据标准化:对性别、肥胖水平、交通方式等字段进行标准化处理
- 衍生字段计算:计算BMI、年龄分组、肥胖等级编码、风险评分
- 质量监控:使用Hadoop Counter统计数据质量指标
3.2 Hive ETL流程
3.2.1 性能优化配置
SET hive.exec.parallel=true;
SET hive.exec.parallel.thread.number=16;
SET hive.auto.convert.join=true;
SET hive.map.aggr=true;
SET hive.groupby.skewindata=true;
SET hive.optimize.skewjoin=true;
SET hive.exec.dynamic.partition=true;
SET hive.exec.dynamic.partition.mode=nonstrict;
SET hive.vectorized.execution.enabled=true;
SET hive.vectorized.execution.reduce.enabled=true;
3.2.2 ODS到DWD层转换
加载性别维度表:
INSERT OVERWRITE TABLE dwd_gender_info
SELECT
DISTINCT gender
FROM ods_obesity_raw
WHERE gender IS NOT NULL
AND gender != '';
加载年龄分组维度表:
INSERT OVERWRITE TABLE dwd_age_group_info
SELECT
age_group,
MIN(age) AS min_age,
MAX(age) AS max_age,
CASE age_group
WHEN 'Adolescent' THEN '青少年(≤18岁)'
WHEN 'Young_Adult' THEN '青年(19-30岁)'
WHEN 'Middle_Aged' THEN '中年(31-50岁)'
WHEN 'Senior' THEN '中老年(51-65岁)'
WHEN 'Elderly' THEN '老年(>65岁)'
END AS description
FROM ods_obesity_raw
WHERE age_group IS NOT NULL
AND age_group != ''
GROUP BY age_group;
加载肥胖风险事实表:
INSERT OVERWRITE TABLE dwd_obesity_risk_fact PARTITION (dt='2026-02-02')
SELECT
id,
gender,
age,
age_group,
height,
weight,
bmi,
CASE
WHEN bmi < 18.5 THEN 'Underweight'
WHEN bmi < 25 THEN 'Normal'
WHEN bmi < 30 THEN 'Overweight'
ELSE 'Obese'
END AS bmi_group,
family_history_with_overweight,
favc,
fcvc,
ncp,
caec,
smoke,
ch2o,
scc,
faf,
tue,
calc,
mtrans,
obesity_level,
obesity_level_code,
risk_score
FROM ods_obesity_raw;
3.2.3 DWS层汇总计算
性别肥胖分析汇总:
INSERT OVERWRITE TABLE dws_gender_obesity PARTITION (dt='2026-02-02')
SELECT
gender,
COUNT(*) AS total_count,
SUM(CASE WHEN obesity_level_code >= 5 THEN 1 ELSE 0 END) AS obesity_count,
SUM(CASE WHEN obesity_level_code IN (3, 4) THEN 1 ELSE 0 END) AS overweight_count,
SUM(CASE WHEN obesity_level_code = 2 THEN 1 ELSE 0 END) AS normal_count,
SUM(CASE WHEN obesity_level_code = 1 THEN 1 ELSE 0 END) AS underweight_count,
ROUND(SUM(CASE WHEN obesity_level_code >= 5 THEN 1 ELSE 0 END) * 100.0 / COUNT(*), 2) AS obesity_rate,
ROUND(SUM(CASE WHEN obesity_level_code IN (3, 4) THEN 1 ELSE 0 END) * 100.0 / COUNT(*), 2) AS overweight_rate,
ROUND(AVG(bmi), 2) AS avg_bmi,
ROUND(AVG(risk_score), 2) AS avg_risk_score
FROM dwd_obesity_risk_fact
WHERE dt='2026-02-02'
GROUP BY gender;
3.2.4 ADS层应用数据
肥胖趋势分析:
INSERT OVERWRITE TABLE ads_obesity_trend_analysis PARTITION (dt='2026-02-02')
SELECT
age_group,
gender,
total_count,
obesity_count,
obesity_rate,
avg_bmi,
avg_risk_score
FROM dws_gender_obesity
WHERE dt='2026-02-02';
实现说明:
- 使用单独的INSERT语句替代UNION ALL
- 支持按日期分区,便于数据管理和查询
- 启用Hive向量化执行,提升查询性能
- 启用并行执行,充分利用集群资源
3.3 后端API实现
3.3.1 Controller层
ObesityAnalysisController.java 提供RESTful API接口:
@RestController
@RequestMapping("/api/obesity")
public class ObesityAnalysisController {
@Resource
private ObesityAnalysisService obesityAnalysisService;
@GetMapping("/trend")
public List<AdsObesityTrendAnalysis> getObesityTrend() {
return obesityAnalysisService.getObesityTrendAnalysis();
}
@GetMapping("/trend/page")
public Object getObesityTrendPage(@RequestParam(defaultValue = "1") int current,
@RequestParam(defaultValue = "10") int size) {
Page<AdsObesityTrendAnalysis> page = new Page<>(current, size);
return obesityAnalysisService.getObesityTrendAnalysisPage(page);
}
@GetMapping("/risk-factors")
public List<AdsObesityRiskFactors> getRiskFactors() {
return obesityAnalysisService.getObesityRiskFactors();
}
@GetMapping("/demographic")
public List<AdsObesityDemographicAnalysis> getDemographicAnalysis() {
return obesityAnalysisService.getDemographicAnalysis();
}
@GetMapping("/lifestyle")
public List<AdsObesityLifestyleAnalysis> getLifestyleAnalysis() {
return obesityAnalysisService.getLifestyleAnalysis();
}
@GetMapping("/comprehensive")
public List<AdsObesityComprehensiveAnalysis> getComprehensiveAnalysis() {
return obesityAnalysisService.getComprehensiveAnalysis();
}
@GetMapping("/raw")
public List<OdsObesityRaw> getRawObesityData() {
return obesityAnalysisService.getRawObesityData();
}
@GetMapping("/health")
public String health() {
return "Obesity Analysis API is healthy!";
}
}
3.3.2 Service层
ObesityAnalysisServiceImpl.java 实现业务逻辑:
@Service
public class ObesityAnalysisServiceImpl implements ObesityAnalysisService {
@Resource
private AdsObesityTrendAnalysisMapper adsObesityTrendAnalysisMapper;
@Resource
private AdsObesityRiskFactorsMapper adsObesityRiskFactorsMapper;
@Resource
private AdsObesityDemographicAnalysisMapper adsObesityDemographicAnalysisMapper;
@Resource
private AdsObesityLifestyleAnalysisMapper adsObesityLifestyleAnalysisMapper;
@Resource
private AdsObesityComprehensiveAnalysisMapper adsObesityComprehensiveAnalysisMapper;
@Resource
private OdsObesityRawMapper odsObesityRawMapper;
@Override
public List<AdsObesityTrendAnalysis> getObesityTrendAnalysis() {
return adsObesityTrendAnalysisMapper.selectList(null);
}
@Override
public Page<AdsObesityTrendAnalysis> getObesityTrendAnalysisPage(Page<AdsObesityTrendAnalysis> page) {
return adsObesityTrendAnalysisMapper.selectPage(page, null);
}
@Override
public List<AdsObesityRiskFactors> getObesityRiskFactors() {
return adsObesityRiskFactorsMapper.selectList(null);
}
@Override
public List<AdsObesityDemographicAnalysis> getDemographicAnalysis() {
return adsObesityDemographicAnalysisMapper.selectList(null);
}
@Override
public List<AdsObesityLifestyleAnalysis> getLifestyleAnalysis() {
return adsObesityLifestyleAnalysisMapper.selectList(null);
}
@Override
public List<AdsObesityComprehensiveAnalysis> getComprehensiveAnalysis() {
return adsObesityComprehensiveAnalysisMapper.selectList(null);
}
@Override
public List<OdsObesityRaw> getRawObesityData() {
return odsObesityRawMapper.selectList(null);
}
@Override
public Page<OdsObesityRaw> getRawObesityDataPage(Page<OdsObesityRaw> page) {
return odsObesityRawMapper.selectPage(page, null);
}
}
3.3.3 Entity层
AdsObesityTrendAnalysis.java 实体类:
@Data
@TableName("ads_obesity_trend_analysis")
public class AdsObesityTrendAnalysis {
@TableId(type = IdType.AUTO)
private Long id;
private String ageGroup;
private String gender;
private Long totalCount;
private Long obesityCount;
private BigDecimal obesityRate;
private BigDecimal avgBmi;
private BigDecimal avgRiskScore;
private String dt;
private LocalDateTime createTime;
private LocalDateTime updateTime;
}
实现说明:
- RESTful API设计:遵循REST规范
- 分页查询:使用MyBatis-Plus的Page对象实现分页
- 依赖注入:使用@Resource注解进行依赖注入
- ORM映射:使用MyBatis-Plus简化数据库操作
3.4 前端可视化实现
3.4.1 Dashboard页面
Dashboard.vue 主页面展示关键指标和图表:
<template>
<div class="dashboard-container">
<el-card class="welcome-card">
<h2>肥胖风险分析系统</h2>
<p>欢迎使用肥胖风险分析系统,这里展示了肥胖趋势、人口统计、生活方式等多维度的分析数据。</p>
</el-card>
<div class="stats-container">
<el-card class="stat-card">
<div class="stat-icon total"><i class="el-icon-user"></i></div>
<div class="stat-info">
<h3>总样本数</h3>
<p class="stat-value">{{ totalSamples }}</p>
</div>
</el-card>
<el-card class="stat-card">
<div class="stat-icon obesity"><i class="el-icon-warning"></i></div>
<div class="stat-info">
<h3>肥胖人数</h3>
<p class="stat-value">{{ obesityCount }}</p>
</div>
</el-card>
<el-card class="stat-card">
<div class="stat-icon rate"><i class="el-icon-data-line"></i></div>
<div class="stat-info">
<h3>肥胖率</h3>
<p class="stat-value">{{ obesityRate }}%</p>
</div>
</el-card>
<el-card class="stat-card">
<div class="stat-icon bmi"><i class="el-icon-s-data"></i></div>
<div class="stat-info">
<h3>平均BMI</h3>
<p class="stat-value">{{ avgBmi }}</p>
</div>
</el-card>
</div>
<div class="charts-container">
<div class="chart-row">
<el-card class="chart-card">
<div slot="header" class="chart-header">
<span>肥胖趋势分析</span>
<el-button type="primary" size="small" @click="navigateTo('/trend-analysis')">查看详情</el-button>
</div>
<div ref="trendChart" class="chart"></div>
</el-card>
<el-card class="chart-card">
<div slot="header" class="chart-header">
<span>风险因素分析</span>
<el-button type="primary" size="small" @click="navigateTo('/risk-factors-analysis')">查看详情</el-button>
</div>
<div ref="riskFactorsChart" class="chart"></div>
</el-card>
</div>
</div>
</div>
</template>
3.4.2 Vuex状态管理
store/index.js 管理全局状态:
export default new Vuex.Store({
state: {
loading: false,
obesityTrendData: [],
riskFactorsData: [],
demographicData: [],
lifestyleData: [],
comprehensiveData: []
},
getters: {
getObesityTrendData: state => state.obesityTrendData,
getRiskFactorsData: state => state.riskFactorsData,
getDemographicData: state => state.demographicData,
getLifestyleData: state => state.lifestyleData,
getComprehensiveData: state => state.comprehensiveData
},
mutations: {
setLoading(state, status) {
state.loading = status
},
setObesityTrendData(state, data) {
state.obesityTrendData = data
},
setRiskFactorsData(state, data) {
state.riskFactorsData = data
},
setDemographicData(state, data) {
state.demographicData = data
},
setLifestyleData(state, data) {
state.lifestyleData = data
},
setComprehensiveData(state, data) {
state.comprehensiveData = data
}
},
actions: {
fetchObesityTrendData({ commit }) {
commit('setLoading', true)
return new Promise((resolve, reject) => {
fetch('http://localhost:8080/api/api/obesity/trend')
.then(response => response.json())
.then(data => {
commit('setObesityTrendData', data)
commit('setLoading', false)
resolve(data)
})
.catch(error => {
commit('setLoading', false)
reject(error)
})
})
}
}
})
实现说明:
- 组件化设计:使用Vue组件化开发
- 状态管理:使用Vuex管理全局状态
- 数据可视化:使用ECharts实现图表展示
- 响应式布局:适配不同屏幕尺寸
四、技术实现要点
4.1 MapReduce实现要点
- HDFS连接检测:检测HDFS可用性,失败时使用本地文件系统
- 数据质量监控:使用Hadoop Counter统计数据质量指标
- 衍生字段计算:在Map阶段计算BMI、风险评分等衍生字段
- 数据标准化:对性别、肥胖水平等字段进行标准化处理
4.2 Hive ETL实现要点
- 移除UNION ALL:使用单独的INSERT语句
- 并行执行:启用并行执行,充分利用集群资源
- 向量化执行:启用Hive向量化执行,提升查询性能
- 动态分区:支持按日期分区,便于数据管理和查询
- 小文件合并:配置小文件合并参数,减少NameNode压力
4.3 后端API实现要点
- RESTful设计:遵循REST规范
- 分页查询:使用MyBatis-Plus的Page对象实现分页
- 依赖注入:使用Spring的依赖注入
- ORM映射:使用MyBatis-Plus简化数据库操作
4.4 前端实现要点
- 组件化开发:使用Vue组件化开发
- 状态管理:使用Vuex管理全局状态
- 数据可视化:使用ECharts实现图表展示
- 响应式布局:适配不同屏幕尺寸
五、项目部署
5.1 环境要求
- JDK 1.8+
- Hadoop 3.x
- Hive 3.x
- MySQL 5.7+
- Node.js 14+
- Maven 3.6+
5.2 部署步骤
- 配置Hadoop环境
- 配置Hive环境
- 编译MapReduce程序
- 执行MapReduce数据清洗
- 执行Hive建表和ETL脚本
- 配置Sqoop数据导出
- 启动Spring Boot后端服务
- 启动Vue前端服务
六、总结
本文介绍了基于Hadoop、Hive、SpringBoot、Vue.js等技术栈实现的肥胖风险数据分析系统。系统实现了以下功能:
- 数据清洗:使用MapReduce对原始数据进行清洗和转换
- 数据仓库:构建ODS/DWD/DWS/ADS四层数据仓库
- 数据分析:实现多维度数据分析
- 数据可视化:使用ECharts实现数据可视化展示
- API服务:提供RESTful API接口
本文提供了主要技术实现的代码示例,供读者参考。
更多推荐


所有评论(0)