电力能耗数据分析系统:Hadoop+Spark+Django全链路实战
做电力能耗数据分析系统绕不开 Hadoop、Spark、Django 这三个词。很多同学一看到“大数据”就发怵其实把它拆开就是一条数据流水线Hadoop 负责海量数据的存储和管理Spark 负责批量计算Django 负责把计算结果变成可视化大屏上的图表和指标。去年我完整跑通了一个带源码、文档、调试记录和可视化大屏的版本过程中踩了不少坑也总结出了一套能直接复用的套路。这篇文章把整个项目的思路、环境配置、数据链路、服务端实现和调试经验都聊一遍无论你是要做课程设计、毕业设计还是接手一个真实的企业能耗分析项目都能从中找到可以落地的东西。1. 项目整体认知与选型思考1.1 这个系统到底解决什么问题电力能耗数据的第一个特点是量大。一个普通园区如果接入了上千只智能电表每 15 分钟采集一条数据一天就是近百万条记录一个月轻松上亿。第二个特点是维度多除了时间还有区域、楼宇、设备类型、电压电流、有功无功等多个字段。第三个特点是分析需求杂既要看实时趋势又要做峰谷统计、区域排名、异常告警。传统方案用一台 MySQL 单库就能做几十万量级的报表但数据量冲到千万级以上后聚合查询会变得很慢一个按楼宇分组的统计可能要跑到几十秒大屏每隔 5 秒刷一次根本扛不住。这个系统的核心目标就是构建一条能够“撑住上亿级数据”的分析链路先让 Hadoop 把原始数据收进来再用 Spark 做分布式清洗和聚合最后把汇总结果落到 MySQL由 Django 对外提供接口和大屏展示。我这里特别想强调一个点不是所有数据都要进 Hadoop。我们要做的是让“最重的计算”跑到分布式集群上然后把“轻量的结果”交给 Web 层。如果一上来就让 Django 直接查 HDFS 或者 Spark 的结果文件开发成本高不说实时性也很难保证。1.2 为什么是 Hadoop Spark Django 这套组合这套组合里面Hadoop 承担的是“仓库”角色。HDFS 存储原始电能数据YARN 负责资源调度Hive 或者 Spark SQL 都可以在上面跑 SQL。Spark 承担的是“加工车间”角色它比 MapReduce 快很多尤其是迭代式计算和交互式查询能把复杂的清洗逻辑和统计分析直接写成 SQL 或 DataFrame 代码适合课程设计和快速开发。Django 承担的是“服务台”角色。它提供 ORM、Admin 后台、模板渲染和成熟生态非常适合快速搭建业务接口和管理界面。有人会问为什么不用 Flask因为这类项目通常要写完整文档、画系统架构图Django 自带 Admin 和管理员模块更容易把“管理后台”这一块说清楚毕业论文和项目评审时也更有说服力。再补充一句如果只用 Spark 做计算不搭 Django数据算完之后摆在哪里是个问题。如果只用 Django没有 Hadoop 和 Spark数据量一大就会变成灾难。三者合在一起各管一段数据从采集到展示的路径才是完整闭环。1.3 适合谁参考这个项目非常适合数据科学与大数据技术专业的学生拿来当毕业设计或者课程设计因为题目里同时覆盖了数据采集、存储、离线计算、接口开发、可视化展示每个环节都能写出东西。对于刚入门大数据开发的人它也是一个很好的“全链路”案例能让你看到 Hadoop、Spark、Django 是怎么配合的而不只是停留在单个组件的小 Demo。如果你是企业里做能耗管理或者智慧园区的人也可以参考这套思路。中小规模的能耗系统不需要太高深的技术栈Hadoop 存原始数据Spark 做清洗和汇总Django 出接口再套一个 ECharts 大屏基本就能满足日常需求。2. 环境搭建与集群配置要点2.1 Hadoop 伪分布式还是集群模式我建议第一阶段先做伪分布式。伪分布式的本质是“一个进程扮演所有角色”HDFS 的 NameNode、DataNodeYARN 的 ResourceManager、NodeManager 都跑在同一台机器上。好处是配置简单、资源占用少适合把整个项目的代码逻辑先跑通。坏处是它不能真正体现 Spark 的分布式优势Executor 还是在本地调度。伪分布式的关键配置集中在三个文件。core-site.xml里要指定默认文件系统地址hdfs-site.xml里要指定副本数和 NameNode 数据目录yarn-site.xml里要指定调度器类型和资源分配。我用的版本是 Hadoop 3.3.4 JDK 8具体配置片段如下!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration!-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name valuefile:///opt/hadoop/data/namenode/value /property property namedfs.datanode.data.dir/name valuefile:///opt/hadoop/data/datanode/value /property /configuration!-- yarn-site.xml -- configuration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property property nameyarn.scheduler.maximum-allocation-mb/name value8192/value /property /configuration伪分布式有两个容易踩的坑。第一个是JAVA_HOME没配Hadoop 启动脚本找不到 JDK直接报错。第二个是不要反复执行hdfs namenode -format格式化一次之后如果 DataNode 目录里的新旧 clusterID 不一致NameNode 很可能启动失败。2.2 三节点 Spark 集群搭建中的坑如果你有条件建议还是搭一个三个节点的集群因为很多问题只有集群环境才会暴露。我的规划是1 台 master2 台 worker。master 节点负责 Spark Master 和调度worker 节点负责执行计算。搭建过程中最容易出问题的不是安装过程而是版本配套。我最终稳定运行的版本组合是Hadoop 3.3.4、Spark 3.2.3、Scala 2.12.15、JDK 8。注意 Spark 3.2 默认支持 Scala 2.12如果你用 Scala 2.13 的编译包OpenCV 或某些第三方库可能会不兼容。spark-env.sh是核心配置我贴一下关键内容export JAVA_HOME/usr/local/jdk1.8 export HADOOP_HOME/opt/hadoop export SPARK_MASTER_HOST192.168.1.100 export SPARK_WORKER_MEMORY4g export SPARK_WORKER_CORES2 export SPARK_HISTORY_OPTS-Dspark.history.fs.logDirectoryhdfs://master:9000/spark-logs另外还要编辑workers文件写上两个 worker 节点的 IP 或主机名192.168.1.101 192.168.1.102启动后访问 Spark Web UI 的 8080 端口正常情况下能看到两个 worker 都在 ALIVE 状态。常见的现象是worker 进程能起来但 Web UI 里显示 DEAD原因是 master 地址写错或者 worker 节点和 master 节点之间的 7077 端口不通。排查时先确认防火墙再确认SPARK_MASTER_HOST是否和当前 master 实际 IP 一致。2.3 Hadoop 与 ZooKeeper 整合实战ZooKeeper 在这个项目里的主要作用是给 Hadoop 提供高可用能力。如果只有单 NameNode一旦机器宕机整个 HDFS 就不可用了。配置 HA 之后两台 NameNode 一主一备通过 ZooKeeper 自动切换。打开core-site.xml加入以下配置configuration property namefs.defaultFS/name valuehdfs://mycluster/value /property property nameha.zookeeper.quorum/name value192.168.1.100:2181,192.168.1.101:2181,192.168.1.102:2181/value /property property namedfs.nameservices/name valuemycluster/value /property property namedfs.ha.namenodes.mycluster/name valuenn1,nn2/value /property /configuration仅配置一行是远远不够的还需要指定每个 NameNode 的 RPC 地址、共享编辑日志目录和自动故障转移开关。实操时我是先在三个节点上安装 ZooKeeper 3.7.1然后配置myid、zoo.cfg再启动zkServer.sh start用zkServer.sh status检查当前节点的角色是 leader 还是 follower。这一步不复杂但很容易忽略ZooKeeper 必须要奇数个节点才能正确选出 leader生产环境最少 3 个。还要注意如果只是做毕业设计其实可以不部署 HA单节点完全能跑。我把 ZooKeeper 放进来是因为真实企业环境里 HBase、Kafka 这些组件都依赖它提前把整合流程走一遍面试时被问到的概率很高。3. 数据链路与核心计算逻辑3.1 电力能耗数据从采集到入库的完整链路我先定义一下原始数据的格式。假设每一条记录包含以下字段时间戳、设备ID、楼宇ID、区域编码、电能量、电压、电流、有功功率、无功功率。在实验环境下我用 Python 写了一个数据生成器随机生成过去 90 天、每天 2 万条记录这样既能模拟真实数据又不会让集群资源爆掉。生成文件之后直接用命令上传到 HDFShdfs dfs -mkdir -p /data/energy/raw hdfs dfs -put energy_data_*.csv /data/energy/raw/接下来是数据清洗。清洗包括几块去掉时间戳为空的行去掉电能量为负数的异常值去掉电压或电流超出合理范围的记录还要把时间戳转成标准格式。我用 Spark SQL 读入 CSV 并做过滤CREATE OR REPLACE TEMP VIEW raw_energy AS SELECT * FROM csv./data/energy/raw/*.csv然后执行清洗逻辑CREATE TABLE cleaned_energy AS SELECT from_unixtime(ts, yyyy-MM-dd HH:mm:ss) AS event_time, device_id, building_id, region_code, electricity_kwh, voltage, current, active_power, reactive_power, substring(from_unixtime(ts, yyyy-MM-dd), 1, 10) AS dt FROM raw_energy WHERE ts IS NOT NULL AND electricity_kwh 0 AND voltage BETWEEN 180 AND 260 AND current BETWEEN 0 AND 100;把清洗后的数据按dt做分区后续按天统计会非常快。注意直接读取 CSV 时Spark 的 schema 推断偶尔会把数值字段识别成字符串建议在读取时显式指定类型否则后面求和会出现全零的情况。3.2 Spark SQL 分析任务设计与执行我实现的几个核心分析需求包括按天统计总用电量、按时段统计峰谷电量、按楼宇和区域排名、识别用电异常。以峰谷电量为例子采用常见的时段划分峰段 08:00-22:00谷段 22:00-次日 08:00。这个划分在代码里直接用 when 条件写SELECT dt, building_id, SUM(CASE WHEN hour(event_time) 8 AND hour(event_time) 22 THEN electricity_kwh ELSE 0 END) AS peak_kwh, SUM(CASE WHEN hour(event_time) 22 OR hour(event_time) 8 THEN electricity_kwh ELSE 0 END) AS valley_kwh FROM cleaned_energy GROUP BY dt, building_id;统计结果一般不会太大因为经过了按楼宇和天聚合。我把结果通过 JDBC 写入 MySQL方便 Django 直接查。写入代码如下df.write.mode(overwrite) .option(truncate, true) .jdbc(jdbc:mysql://192.168.1.10:3306/energy_web, daily_building_energy, props)提交任务的命令也很重要不要直接在终端跑spark-shell要用spark-submitspark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 2g \ --executor-memory 4g \ --num-executors 3 \ --executor-cores 2 \ --class com.energy.analyzer.DailyAnalyzer \ energy_analyzer.jar如果你用 Python也可以类似提交 PySpark 脚本。这里的核心经验是executor-memory不能设置得比 YARN 容器最大值还大否则任务会一直卡在 ACCEPTED 状态。3.3 数据倾斜和内存调优经验按楼宇聚合时最容易出现数据倾斜。因为某些大楼的能耗记录数量可能是其他楼宇的几十倍分配到某个 Executor 的 Task 会处理超大数据量导致整个 Stage 被拖住。我遇到过一个很明显的案例一个 Stage 有 200 个 Task199 个几秒钟就做完了剩下 1 个跑了 20 分钟还没结束就是典型的倾斜。解决办法是加盐。先把 key 加上一个随机前缀让数据分散到多个分区完成第一轮聚合后再去掉前缀做第二轮聚合。用 PySpark 写的话大致是from pyspark.sql.functions import rand, split, col salted df.withColumn(salt, (rand() * 10).cast(int)) \ .withColumn(salted_building, concat(col(building_id), lit(_), col(salt))) part1 salted.groupBy(salted_building).sum(electricity_kwh) \ .withColumn(building_id, split(col(salted_building), _)[0]) result part1.groupBy(building_id).sum(sum(electricity_kwh))内存调优方面我建议关注三个点。第一是spark.sql.shuffle.partitions默认 200如果数据量不大可以调低到 64减少小文件数量。第二是启用 Kryo 序列化能明显降低内存占用。第三是留出足够的spark.yarn.executor.memoryOverhead默认只占 executor 内存的 10%在做大聚合时容易 OOM我通常调到 512M 或 1G。4. Django 服务端与可视化大屏实现4.1 Django 如何读取 Spark 计算结果Spark 跑完的结果已经写入 MySQLDjango 这边只需要把它当作普通的业务表来查。这样设计的好处是 Django 不需要写复杂的分布式计算代码也不需要在接口请求时临时去调 Spark前端请求的响应时间能控制在几百毫秒以内。模型定义大概是这样的from django.db import models class DailyBuildingEnergy(models.Model): dt models.DateField() building_id models.CharField(max_length32) peak_kwh models.FloatField() valley_kwh models.FloatField() total_kwh models.FloatField() class Meta: db_table daily_building_energy在视图里用 ORM 聚合能少写很多原生 SQL。例如要查询最近 7 天总耗电量可以这样from django.db.models import Sum from django.http import JsonResponse from .models import DailyBuildingEnergy def energy_trend(request): rows DailyBuildingEnergy.objects.values(dt) \ .annotate(totalSum(total_kwh)) \ .order_by(dt)[:7] data [{date: item[dt], total: item[total]} for item in rows] return JsonResponse({code: 0, data: data})ORM 并不是万能的如果聚合条件特别复杂也可以使用raw()方法写原生 SQL。但要注意原生 SQL 里的表名必须和 MySQL 中的实际表名一致否则会报Table not found。4.2 可视化大屏的数据接口设计大屏前端我用的是 ECharts 普通 HTML 页面没有引入 Vue 或 React。因为大屏的核心是快速展示不需要复杂交互。后台只需要提供几个 JSON 接口格式统一为{code: 0, msg: success, data: {...}}。我设计了三个核心接口/api/energy/overview当天总用电量、峰电、谷电、环比去重率。/api/energy/trend近 7 天或近 30 天的用电趋势。/api/energy/rank按区域或楼宇的用电排行。接口实现里要记得加缓存。大屏一般是 5 秒轮询一次但像总用电量这种数据完全没必要每次都算一遍用django.core.cache.cache_page缓存 30 秒就够了from django.views.decorators.cache import cache_page from django.views.decorators.http import require_GET require_GET cache_page(30) def energy_overview(request): ...如果 Django 服务和前端页面不在同一个域名下记得安装django-cors-headers并在 settings 里配置白名单。不然浏览器会直接拦截请求大屏上所有请求都是红叉。4.3 Django 执行查询-删除对象时容易忽略的问题这部分是很多初学者容易忽略的。Django QuerySet 是惰性的你在filter()之后不会立刻去数据库查只有迭代、取切片、调用list()或bool()时才会真正执行查询。如果你写了一段代码filter 之后又多次遍历同一个 QuerySet每次遍历都会重新触发 SQL性能会很差。正确做法是先转成列表qs DailyBuildingEnergy.objects.filter(dt2024-05-01) energy_list list(qs)删除对象也有讲究。千万不要在 Python 层循环删除单个对象比如for obj in qs: obj.delete()这会造成大量数据库操作。正确姿势是用 QuerySet 的delete()DailyBuildingEnergy.objects.filter(dt2024-05-01).delete()delete()的返回值是一个元组第一个元素是受影响总行数第二个元素是一个字典包含每个模型被删除的具体数量。如果模型有关联外键并且设置了on_deletemodels.CASCADE删除主表数据时关联表也会被级联删除。这个特性在某些场景下很危险比如你想清空日志表结果把关联的告警表也删了。所以每次delete()之前先把返回的字典打印出来确认一下。批量删除大量数据时单条 SQL 可能占用过长时间。经验做法是分批删除每次只删一千条while DailyBuildingEnergy.objects.filter(dt2024-05-01).exists(): DailyBuildingEnergy.objects.filter(dt2024-05-01)[:1000].delete() time.sleep(1)5. 调试、部署与常见问题速查5.1 联调中最容易出现的三类问题第一类是端口不通。HDFS 默认 RPC 端口是 9000Spark Master 是 8080Spark 提交任务是 7077MySQL 是 3306Django 开发服务默认是 8000。联调时经常出现 Web 页面访问不了但进程还在就是因为防火墙没放行。排查命令很简单ss -tlnp | grep 8080 curl http://192.168.1.100:8080第二类是版本冲突。最典型的是 Hadoop 2.x 和 Spark 3.x 一起用Spark 在读取 HDFS 时会找不到某些类是常见报错。建议所有组件版本都选较新的稳定组合比如 Hadoop 3.3.x 配 Spark 3.2.x并且保证集群所有节点的 JDK 版本一致。第三类是权限问题。HDFS 上的数据目录默认权限是 700如果其他用户执行 Spark 任务时没有权限读取会报Permission denied。可以直接改成 755hdfs dfs -chmod -R 755 /data/energyMySQL 那边的坑则是用户权限。Spark 写入时用的账号可能只赋予了localhost权限而 Spark 任务跑在远程节点会报Access denied for user。需要在 MySQL 里给账号添加远程访问权限并指定具体网段。5.2 环境重启顺序和资源分配建议开发期间免不了反复重启集群。我建议严格按照以下顺序启动先启动 ZooKeeper然后启动 Hadoop 的 HDFS 和 YARN最后启动 Spark。不要在 HDFS 还没进入安全模式之前就启动 Spark否则 Spark 作业初始化会失败。关闭顺序反过来先停 Spark再停 YARN 和 HDFS最后停 ZooKeeper。我用的是一台 8核 16GB 的服务器加两台 4核 8GB 的节点资源分配大致如下服务内存CPU 核数Hadoop NameNode ZooKeeper master2G1Hadoop DataNode ZooKeeper server1G1Spark Master1G1Spark Worker4G2MySQL Django2G1系统预留1G1资源不足时宁可减少 Executor 数量也不要让单个 Executor 的内存小于 2G否则大批量聚合时基本都会 OOM。5.3 实用排查命令与经验遇到进程挂掉先看jps输出只要能看到 NameNode、DataNode、ResourceManager、NodeManager基本说明 Hadoop 这边没问题。如果某个进程缺失去对应的日志目录找*.log文件。Hadoop 日志在$HADOOP_HOME/logsSpark 日志在$SPARK_HOME/logsDjango 调试日志建议配一个单独的文件输出不要只打在终端上。Spark 任务执行失败时第一反应是去 Spark UI 看 Job / Stage 的日志摘要而不是盲猜。点开失败 Stage 的 Executor 日志通常能直接看到java.lang.OutOfMemoryError或者org.apache.spark.shuffle.FetchFailedException。后者大概率是 Executor 掉线或网络抖动可以重试一次如果反复出现再排查内存和网络。MySQL 慢查询也是排查重点。在 my.cnf 开启慢查询日志slow_query_log ON long_query_time 1如果发现某个 Django 接口访问特别慢先看是不是命中了缓存再看看是不是 ORM 产生了 N1 查询。大屏页面需要一次性返回多个图表数据时尽量一个接口打包所有数据不要每个图表单独请求一次。最后分享一个小技巧用tmux跑spark-submit的长任务。我曾经因为 SSH 断开会话导致一个跑了半小时的清洗任务直接丢失数据没写到 MySQL大屏上全是 0。从那以后我学乖了先tmux new -s spark_task新开会话再在会话里提交任务这样即使断开连接任务也能继续跑完。如果你也要做这个题我的建议是先跑通一条最简链路手动生成 1 万条数据Hadoop 上传Spark SQL 统计Django 输出接口ECharts 画一张趋势图。这条链路通了之后再做峰谷统计、排名和异常告警整个项目的完成度会提升非常快。