十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Hadoop+Spring Boot:电力生产数据分析系统实战

Hadoop+Spring Boot:电力生产数据分析系统实战 简介基于Hadoop大数据与Spring Boot的电力生产数据分析系统源码项目面向计算机相关专业学生、毕业设计者与入门开发者覆盖电力数据从HDFS存储、PySpark预处理分析到Web可视化展示的完整业务链路适合作为毕业设计、课程设计或项目初期演示。资源共包含三百六十九个文件压缩包大小仅9.6MB以Java后端源码、Vue前端页面、Python数据分析脚本、XML配置文件和SQL初始化脚本为主并附有项目截图、使用说明与搭建文档可支撑环境配置、代码调试和功能演示。当前已有一百六十六人学习或下载属于典型的毕设高分项目。读者可从中获得完整的数据表结构、分析算法实现、Web交互界面及部署配置思路还可对照搭建文档完成环境部署与代码调试在此基础上修改扩展用于课程设计、项目立项演示或大数据技术实践。1. 电力生产数据分析系统从 Hadoop 到 Spring Boot 的一条真实数据链路做电力数据的人都有个共同的痛数据量一上来Excel 直接卡死MySQL 单表几千万行之后连 count 都要等半天。这个基于 Hadoop Spring Boot 的电力生产数据分析系统把存储层放在 HDFS用 PySpark 做预处理再由 Spring Boot 提供 REST 接口配合 Vue 页面做可视化——整套链路是完整跑通的不是只有前端页面或者只有后端接口的半吊子毕设。从 CSV 原始文件到最终图表展示每一步都有对应代码和验证方式。适合两类人一类是准备毕设答辩、需要一套能讲清楚架构的在校生另一类是工作中刚接触大数据生态、想抄一条现成数据链路做参考的 Java 工程师。这里先给结论项目能上 96 分不是因为算法多花哨而是因为每个环节都是可以演示、可以讲原理的。你拿到手之后不需要改任何核心代码按照 README 把环境配好就能跑起来。下面从选型开始把这套系统的骨架一层层拆开。2. 为什么是 Hadoop Spark Spring Boot这套组合的边界与取舍2.1 存储层选型HDFS 和 YARN 在电力数据场景里承担什么角色电力生产数据最大的特点是有时序性——发电量、负荷、电压、温度这些指标按时间连续产生一天就能积累上千万条记录。这个项目里的PowerData.csv和dataPower.csv就是典型的按站点、按时间组织的生产数据。HDFS 在这里不是用来替代 MySQL而是充当原始数据仓库CSV 文件扔进去YARN 负责给后续的 PySpark 任务分配计算资源。以我拆过的大数据毕设项目来看很多同学会把 HDFS 和 MySQL 的职责搞混直接在 MySQL 里建一个大宽表硬扛。但那不是大数据项目的思路。HDFS 的优势是顺序读写吞吐量大几十个 GB 的 CSV 文件能稳定写进去YARN 则解决多个计算任务怎么分资源的问题。比如你要对三年的电力数据进行月度聚合YARN 会把 Task 分发到不同节点上并行跑而不是靠单机 CPU 硬算。这个项目的架构里HDFS 上存的是原始 CSV经过 PySpark 清洗之后的结果会落到 MySQL通过 Druid 连接池管理Spring Boot 负责对外提供接口。所以 HDFS 是数据湖MySQL 是服务库职责清晰。2.2 计算层选型为什么数据预处理非要用 PySpark 而不是纯 Java先看一段这个项目里最有代表性的 PySpark 数据清洗代码逻辑很简单但足够典型from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, avg, date_format # 初始化 SparkSessionlocal 模式下用 4 个线程模拟集群 spark SparkSession.builder \ .appName(power_data_etl) \ .master(local[4]) \ .enableHiveSupport() \ .getOrCreate() # 读取 HDFS 上的原始 CSV自动推断 Schema df spark.read \ .option(header, true) \ .option(inferSchema, true) \ .csv(hdfs://localhost:9000/user/hadoop/power/PowerData.csv) # 清洗剔除功率为负的异常记录站点 ID 空值的用 unknown 填充 df_clean df.filter(col(active_power) 0) \ .fillna({station_id: unknown, voltage: 0.0}) # 按站点和小时聚合得到每小时平均发电功率 df_hourly df_clean.groupBy(station_id, date_format(collect_time, yyyy-MM-dd HH)) \ .agg(avg(active_power).alias(avg_power)) \ .withColumnRenamed(date_format(collect_time, yyyy-MM-dd HH), hour_slot) # 写出到 MySQL使用 overwrite 模式避免重复数据 df_hourly.write \ .mode(overwrite) \ .jdbc(urljdbc:mysql://localhost:3306/power_db?useUnicodetruecharacterEncodingutf8, tablehourly_power_stat, properties{user: root, password: 123456, driver: com.mysql.cj.jdbc.Driver}) spark.stop()这段代码里有几个关键参数需要注意。master(local[4])表示在本地用 4 个线程模拟分布式计算如果你机器性能还可以改成local[*]会自动按 CPU 核心数分配enableHiveSupport()意味着你需要在 Hadoop 环境里配好 Hive 元数据服务如果不需要 Hive 表支持这行可以去掉否则会报元数据连接错误。inferSchema(true)会自动识别时间戳和数值类型但如果 CSV 里有脏格式数据建议改成手动指定 Schema保证生产环境的稳定性。聚合逻辑里groupBy(station_id, date_format(...))是我见过最常见的按小时分桶统计输出结果直接喂给 Spring Boot 的查询接口。这里必须说明数据清洗放到 PySpark 做而不是用 Java 在业务层写核心原因是处理效率。几千万条 CSV 如果用 JDBC 逐行读再插 MySQL至少跑半小时PySpark 分布式读 聚合 批量写入两三分钟就完成而且代码量少一半。2.3 服务层选型Spring Boot MyBatis 在这套系统里到底管什么Hadoop 生态里的组件有十几种HBase、Hive、Kafka 各有各的用处但这个项目只选了 Spring Boot MyBatis 做服务端。因为业务场景不是实时流处理而是查统计结果 看趋势图。Spring Boot 在这里的角色是聚合层把 MySQL 里 PySpark 算好的汇总数据包装成 JSON 接口MyBatis 负责和 MySQL 交互。!-- pom.xml 关键依赖 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.mybatis.spring.boot/groupId artifactIdmybatis-spring-boot-starter/artifactId version2.1.4/version /dependency dependency groupIdcom.alibaba/groupId artifactIddruid-spring-boot-starter/artifactId version1.1.21/version /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId scoperuntime/scope /dependencyDruid 连接池在这里不是凑数——PySpark 批量写入 MySQL 时连接数很容易打满Druid 的maxActive和initialSize参数直接决定能不能扛住写入峰值。我一般建议maxActive设为写入并发数的两倍项目里默认配置可以先用如果你的数据量翻倍了再调。MyBatis 在这个项目里管理的表主要是hourly_power_stat这类聚合结果表。你可能会问既然数据已经在 PySpark 里算完了为什么不用 JPA 直接 CRUD答案是 MyBatis 对复杂 SQL 掌控更精确尤其是按时间范围、按多个站点做条件查询时XML 里写 SQL 比派生方法名直观得多。详情见下一章的数据接口实现。3. 数据链路落地从 CSV 上架 HDFS 到 Vue 页面展示的完整步骤3.1 环境准备Hadoop 伪分布式 PySpark 的安装验证在跑通代码之前先确认环境是健康的。Hadoop 伪分布式是毕设最常见的部署方式不需要三台机器一台就行。整个安装过程网上教程很多这里只讲验证步骤和两个坑。# 1. 启动 HDFS 和 YARN start-dfs.sh start-yarn.sh # 2. 检查进程是否齐全NameNode、DataNode、ResourceManager、NodeManager jps # 3. 查看 HDFS 根目录确认文件系统可用 hadoop fs -ls / # 4. 把项目里的原始 CSV 上传到 HDFS hadoop fs -mkdir -p /user/hadoop/power hadoop fs -put PowerData.csv /user/hadoop/power/ hadoop fs -put dataPower.csv /user/hadoop/power/ # 5. 验证文件块分布500MB 以下的文件默认 128MB block会分成多块 hadoop fs -stat %o /user/hadoop/power/PowerData.csv第五步的%o参数可以查看块大小如果你看到 134217728128MB说明 HDFS 配置正常。很多新手在这步会翻车原因是 HDFS 的默认 block 大小配置没生效——检查hdfs-site.xml里的dfs.blocksize是否写对位置注意单位是字节。PySpark 的安装要注意版本匹配Spark 3.x 对应的 PySpark 是 3.x不要装成 2.x 再硬跑pyspark命令。验证方式很简单python3 -c from pyspark.sql import SparkSession; print(SparkSession.builder.getOrCreate().version)能打印出版本号就说明 PySpark 基础环境没问题。如果报 Java 版本错误去确认JAVA_HOME指向的是 JDK 8 而不是 JDK 11——Spark 3.x 对 Java 版本有严格的兼容范围。3.2 数据清洗脚本改造把通用 PySpark 脚本适配到自己的 CSV 字段项目自带的 PySpark 脚本字段名是固定的比如station_id、active_power、collect_time但你的 CSV 列名可能不一样。最常见的场景是字段名带中文或者空格直接跑脚本会报AnalysisException。我的习惯是先用 PySpark 读入数据后打印 Schema确认字段名再做映射# 打印 DataFrame 的字段名和类型 df.printSchema() # 结果示例root |-- 站点编号: string |-- 有功功率: double |-- 采集时间: timestamp看到实际字段名之后把代码里所有col(station_id)改成col(站点编号)即可。注意 PySpark 的col()接受字符串参数中文可以直接传不用额外处理编码。再补充一个聚合粒度的细节如果 CSV 里的时间字段是字符串date_format可能解析失败。处理方式是在读取时用to_timestamp显式转换from pyspark.sql.functions import to_timestamp df df.withColumn(collect_time, to_timestamp(col(collect_time), yyyy-MM-dd HH:mm:ss))to_timestamp的第二个参数是格式模板yyyy是小写HH是 24 小时制——写错大小写会得到全 null这是最常见的翻车点。3.3 Spring Boot 数据接口实现MyBatis 映射与 Druid 连接池配置数据算好了接下来要让 Spring Boot 能查出这些结果。先看application.yml的关键配置spring: datasource: type: com.alibaba.druid.pool.DruidDataSource url: jdbc:mysql://localhost:3306/power_db?useUnicodetruecharacterEncodingutf8useSSLfalseserverTimezoneAsia/Shanghai username: root password: 123456 driver-class-name: com.mysql.cj.jdbc.Driver druid: initial-size: 5 min-idle: 5 max-active: 20 max-wait: 60000 validation-query: SELECT 1 test-while-idle: true mybatis: mapper-locations: classpath:mapper/*.xml type-aliases-package: com.power.entity这里的serverTimezoneAsia/Shanghai一定要加上否则连接 MySQL 8.x 报时区错误。Druid 的validation-query配置成SELECT 1是常规操作用于启动时预检连接别省。对应的 MyBatis Mapper 接口和 XML 如下public interface HourlyPowerMapper { // 按站点和日期范围查询小时级平均功率 ListHourlyPowerStat selectByStationAndTime(Param(stationId) String stationId, Param(startTime) String startTime, Param(endTime) String endTime); }select idselectByStationAndTime resultTypecom.power.entity.HourlyPowerStat SELECT station_id AS stationId, hour_slot AS hourSlot, avg_power AS avgPower FROM hourly_power_stat WHERE station_id #{stationId} AND hour_slot BETWEEN #{startTime} AND #{endTime} ORDER BY hour_slot ASC /select注意查询字段用了别名AS stationId映射到 Java 驼峰属性这是 MyBatis 最常见的坑之一——如果不做别名映射查出来的对象属性全是 null。也可以在application.yml里开启map-underscore-to-camel-case: true但显式别名更直观出问题好排查。Controller 层就很简单了RestController RequestMapping(/api/power) public class PowerController { Autowired private HourlyPowerMapper hourlyPowerMapper; GetMapping(/trend) public Result getPowerTrend(RequestParam String stationId, RequestParam String startTime, RequestParam String endTime) { ListHourlyPowerStat list hourlyPowerMapper.selectByStationAndTime(stationId, startTime, endTime); return Result.success(list); } }接口地址是GET /api/power/trend?stationIdxxxstartTime2024-01-01 00:00endTime2024-01-01 23:00Vue 前端通过 axios 调用这个接口渲染 ECharts 折线图。整套链路从 HDFS 到页面到这里就通了。3.4 前端联调Vue 项目启动与接口代理配置项目前端在static目录下用原生 Vue 或者构建后的 dist 文件都可以直接跑。如果是开发模式注意config/index.js里的 proxyTable 配置——把/api前缀代理到http://localhost:8080避免跨域问题proxyTable: { /api: { target: http://localhost:8080, changeOrigin: true, pathRewrite: { ^/api: /api } } }3.5 常见踩坑启动顺序和数据目录权限整套项目跑不起来的头号原因是启动顺序不对。HadoopHDFS YARN必须先起然后跑 PySpark 脚本写数据到 MySQL最后才是 Spring Boot。如果 MySQL 里还没有聚合表数据后端接口查出来是空数组前端图表自然空白这不是代码 bug是数据流没走完。第二个坑是 HDFS 目录权限。如果你用hadoop用户启动服务却用root用户上传 CSV会报Permission denied。解决方式hadoop fs -chmod -R 777 /user/hadoop/power生产环境不建议这么干但毕设和本地复现直接这样最快。4. 避坑与排查五个真实翻车现场4.1 lombok 插件缺失导致编译失败现象IDEA 编译报红提示找不到getUser()、setStationId()之类的方法但代码看起来没有问题。原因项目用了 Lombok 注解IDEA 默认不带该插件只装了 JDK 是编不过的。解决打开 IDEA 的Settings - Plugins搜索 Lombok 安装并重启 IDEA。注意annotation processing也要打开具体路径是Settings - Build, Execution, Deployment - Compiler - Annotation Processors勾选Enable annotation processing。4.2 PySpark 读取 CSV 中文乱码现象清洗脚本跑完写入 MySQL 的中文全部变成问号。原因HDFS 里文件的编码不是 UTF-8或者 PySpark 读取时没指定编码。解决读取时显式加.option(encoding, UTF-8)同时确认原始 CSV 文件用记事本另存为 UTF-8 编码。还有一个隐蔽原因——MySQL 表字段的charset是latin1建表时务必用utf8mb4否则 Java 端写入时做了编码转换读出来还是乱。4.3 MySQL 连接超时导致 Druid 连接池报CommunicationsException现象Spring Boot 启动正常但长时间不访问后再请求接口第一次请求必然报连接超时错误。原因MySQL 默认wait_timeout是 8 小时连接池里的空闲连接被服务端断开但 Druid 不知道。解决在application.yml里配置test-while-idle: true和time-between-eviction-runs-millis: 60000每 60 秒检测一次空闲连接同时把 MySQL 的wait_timeout调大或者设成 0无限制。这个坑在毕设答辩演示时最容易翻车——演示前 5 分钟没操作一打开页面就报错。4.4 项目启动时 MyBatis 报Invalid bound statement (not found)现象调用 Mapper 接口方法时报错提示找不到对应的 SQL 语句。原因mapper-locations路径配错了或者 XML 文件没编译到target/classes目录下。解决先检查target/classes/mapper/下有没有对应的 XML 文件没有的话说明 IDEA 没有把src/main/resources下的 XML 复制过去。在pom.xml的build节点加一段配置把src/main/java下的 XML 也纳入资源扫描resources resource directorysrc/main/resources/directory filteringtrue/filtering /resource resource directorysrc/main/java/directory includes include**/*.xml/include /includes /resource /resources4.5 前端接口成功但图表不显示现象浏览器 Network 面板能看到接口返回 JSON 数据但页面图表空白。原因ECharts 初始化时series.data的数据格式和接口返回的不匹配。比如后端返回的是avgPower: null前端没做空值过滤。解决在 Vue 的方法里对数据做一层清洗const trendData res.data.data.filter(item item.avgPower ! null)还有一种情况是时间字段的格式——ECharts 的xAxis如果是category类型时间字符串必须连续如果数据库里有小时缺口图表会出现断线。建议聚合时按小时补零或者在axisLabel里设置interval: 0强制显示所有刻度。5. 从毕设到真实业务把项目改造成可运维系统的进阶方向5.1 调度层面用 Azkaban 替代手动跑 PySpark 脚本现阶段每次更新数据都要手动执行一次 PySpark 脚本这在毕设演示里没问题但真实业务场景必须定时调度。把清洗脚本打包后用 Azkaban 配置一个 workflow每天凌晨三点执行 ETL、八点刷新聚合表。核心是在脚本里支持传参--start-date和--end-date增量处理前一天的数据而不是每次全量覆盖。对于刚接触大数据的同行如果不想引入 Azkaban用 Cron Shell 也能凑合但调度历史和失败重试就差远了。5.2 数据一致性聚合表如何应对历史数据修正电力生产数据有一个特点当天采集的数据凌晨会修正比如电表延迟上报所以昨天的聚合结果不能只算一次。我的做法是设计hourly_power_stat表时加一个data_version字段每次 ETL 写入时固定一个版本号。查询时加WHERE data_version (SELECT MAX(data_version) ...)子查询保证前端永远拿到的是一版完整的数据而不是混合状态。这个细节加在毕设答辩里评审会认为你考虑到了生产环境的数据治理问题而不是只会跑通 Demo。5.3 监控与告警给服务加一层体检Spring Boot Actuator 配一个健康检查端点Hadoop 侧用hadoop dfsadmin -report看 DataNode 状态。实际运维中最有效的两个指标是MySQL 里hourly_power_stat表的最大时间戳如果停留在昨天说明 ETL 挂了和 Druid 连接池活跃连接数如果持续超过maxActive的 80%说明查询 SQL 有问题。5.4 优化查询分页插件与索引设计MyBatis 分页插件PageHelper在这个项目中可以有更好的应用。hourly_power_stat表随着运行时间增长累积的数据量会越来越大。为了确保查询性能需要在hour_slot和station_id上建立联合索引。同时接口层面加上分页参数pageNum和pageSize有利于在大规模数据下进一步优化查询效率。从那以后我每次搭类似的项目都会强制走一遍完整的链路验证CSV 能查、HDFS 能读、PySpark 能算、MySQL 能写、Spring Boot 能出接口、页面能渲染六步少一步都不算完成。这个系统的价值不在于代码量多少而在于给你一条可以照着推演的路径。如果你正准备做大数据方向的毕设或者课设从这套代码里挑一条链路拆开研究比你自己从头搭要省一半时间。希望帮到你。本文还有配套的精品资源点击获取
返回列表