ARTICLE DETAIL

资讯详情

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

基于Hadoop+Spark+Django的高血压人群分析可视化系统实战

基于Hadoop+Spark+Django的高血压人群分析可视化系统实战 今年做大数据的毕业设计很多人绕不开医疗健康这个方向。你要是搜过高血压人群分析相关的开源项目大概率会看到一套技术组合Hadoop管历史数据存储、Spark做离线清洗和统计分析、Django出接口、最后用可视化大屏把分析结果铺开展示。这套系统就是按这个思路做的——把“高血压人群分布在哪、哪些年龄段风险高、不同地区的患病率差异如何”这类问题从原始数据一路处理到大屏上的动态图表整个链路完整跑通。这套项目比较适合两类人一类是正在做大数据方向毕业设计、需要一套能讲清楚“数据从哪来、怎么算、怎么展示”的完整案例的同学另一类是刚接触大数据开发、想搞清楚Hadoop、Spark和Web框架之间到底怎么协作的初学者。它的价值不在于某个单点技术用得多深而在于给你一条可以直接抄作业的技术链路HDFS存数据、Spark跑分析、结果入库MySQL、Django做查询接口、前端大屏负责可视化每个环节替换成自己的业务场景就能复用。1. 项目定位与技术架构拆解1.1 高血压人群分析到底在分析什么先说业务侧。高血压分析不是简单统计一个“多少人血压高”而是要做人群画像和风险分层。实际项目里常见的数据来源于体检记录、门诊病历、社区慢病随访表至少包含年龄、性别、收缩压、舒张压、地区、是否吸烟饮酒、体重指数BMI这些字段。高血压的判定标准很好记非同日三次测量收缩压≥140mmHg或舒张压≥90mmHg就属于高血压。但分析系统不能只做一个阈值判断还要做分层——正常高值120-139/80-89、高血压1级140-159/90-99、2级160-179/100-109、3级≥180/≥110。分级之后再叠加年龄、地区、生活习惯这些维度才能看出高危人群的特征。比如“华东地区60-70岁男性、BMI超28、有饮酒史”这样的画像才是对公共卫生管理真正有用的结论。所以这套系统处理的不是单一查询而是一组分析需求不同年龄段患病率趋势、各省份高血压分布、性别差异、血压分级构成比、异常血压人群的区域排名。这些需求拆分到技术实现上就是数据清洗、多维度聚合、统计计算和可视化报表。1.2 为什么是HadoopSparkDjango的组合选这套组合本质上是让每个组件干自己最擅长的事。HDFS用来存原始数据。体检数据往往是批量导入的几百万行CSV是常态HDFS能低成本地把这些文件分散存储在多台机器上而且天然支持副本冗余不用担心单点故障把数据弄丢。Spark用来做计算。它的核心优势是内存计算特别适合多步骤的迭代式数据处理。分析高血压人群要做清洗、过滤、分组聚合、多表关联如果用MapReduce硬写每个步骤都要落盘跑一遍全量数据要等很久。Spark把这些中间结果尽量放在内存里同样是几百万条记录处理速度能快一个数量级。Django用来做查询和展示层。分析结果最终要给人看Django的ORM和DRFDjango Rest Framework写查询接口非常快而且后台管理、用户认证这些现成的功能可以直接用不用从头造轮子。可视化大屏则是数据价值的出口。没有大屏分析结果就是数据库里的几张表领导看不懂、答辩讲不清有了大屏全省高血压分布地图、年龄趋势折线、危险因素排名一眼就能看明白。这套组合也踩过不少坑最典型的就是“数据量并不大却硬上Hadoop”。如果你只有几万条记录单机MySQL加Python脚本完全够用根本不需要分布式。但这个项目本身是教学性质的目的是演示大数据全链路所以Hadoop、Spark是必选项。1.3 系统模块划分与数据流转整个系统按数据流向可以分成四层存储层HDFS存放原始CSV和预处理中间结果MySQL存放Spark算出来的聚合统计表。计算层Spark读取HDFS数据完成清洗、血压分级计算、多维度聚合最后写回MySQL。服务层Django提供REST接口查询MySQL中的统计结果Django Channels提供WebSocket推送。展示层Vue或原生HTMLECharts搭的大屏页面定时拉取接口数据或者等WebSocket推送后自动刷新。数据流转的链路是原始体检数据 → 上传到HDFS → Spark作业读取并清洗 → 统计结果写入MySQL → Django接口查询MySQL → 前端大屏渲染图表。这条链路最关键的设计决策是“计算与展示解耦”。大屏和Django永远不直接查HDFS也不直接触发Spark任务而是只查MySQL里的统计结果。这样前端响应快Spark任务又可以离线批量跑互不干扰。实际项目里我把Spark跑批和Web服务放在两台机器上前端访问高峰期完全不影响后台计算。2. Hadoop与Spark数据链路搭建实战2.1 Hadoop环境从零搭建伪分布式起步很多人在Hadoop环境这一步就卡住了。我的建议很简单先搭伪分布式单机环境把业务逻辑跑通再考虑集群。伪分布式的意义是让你用最小成本把HDFS、YARN、Spark之间的协作关系搞清楚排错的时候不用同时面对十几台机器。Hadoop 3.x的伪分布式安装步骤并不复杂。先确保装了JDK8然后下载Hadoop二进制包解压配置环境变量export HADOOP_HOME/opt/hadoop export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin接着改配置文件。core-site.xml指定NameNode地址configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configurationhdfs-site.xml设置副本数为1伪分布式只有一台机器副本数设成3会报错改成1才正常property namedfs.replication/name value1/value /propertyyarn-site.xml里配置资源管理mapred-site.xml指定用YARN调度。第一次启动前必须格式化NameNodehdfs namenode -format start-dfs.sh start-yarn.sh jpsjps能看到NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager这几个进程就说明启动成功。很多人会漏掉SecondaryNameNode的端口配置导致它起不来我习惯直接在hdfs-site.xml里显式配置property namedfs.namenode.secondary.http-address/name valuelocalhost:9868/value /property伪分布式搭建时最容易犯的错是反复执行namenode -format。格式化一次就够了格式化两次以上会导致NameNode和DataNode的集群ID不一致DataNode起不来。如果遇到这个问题最干脆的解决办法是删除data和name目录、清空tmp再重新格式化。2.2 从伪分布式到集群Hadoop与Zookeeper整合业务数据量大了之后单点NameNode就成了瓶颈。NameNode挂掉意味着整个HDFS不可用所以生产环境必须做高可用。Hadoop HA的标配就是引入Zookeeper由它来协调Active和Standby两个NameNode的自动切换。Zookeeper在整套系统里承担两个职责一是帮NameNode做自动故障转移Active节点挂了之后Standby节点自动顶上二是做YARN ResourceManager的高可用协调。搭建时通常准备三台Zookeeper节点大小是奇数因为Zookeeper选举需要多数派同意。整合的核心操作包括配置journalnode共享存储、配置dfs.nameservices和dfs.ha.namenodes、配置zkfcZookeeper Failover Controller最后在Zookeeper里创建namenode的HA状态节点。启动顺序有讲究先启动Zookeeper集群再启动journalnode然后格式化并启动NameNode最后启动zkfc。顺序错了zkfc就会报连接超时。Hadoop和Zookeeper整合的好处是让整个系统从“能跑”变成“扛得住”。我在课程设计答辩里讲这套HA机制时基本上不需要额外准备太多内容因为Zookeeper本身就是一个独立的考点自动故障切换、分布式协调、选举机制都是可以展开讲的话题。2.3 Spark集群部署与作业提交流程Spark部署我推荐Standalone模式比YARN模式好上手也更容易定位问题。下载与Hadoop版本兼容的Spark包后配置spark-env.shexport JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 export SPARK_MASTER_HOSTnode01 export SPARK_WORKER_CORES4 export SPARK_WORKER_MEMORY8gMaster和Worker都在这台机器上就构成了单机版Spark集群。启动后用浏览器访问8080端口能看到Worker的资源信息。如果之后要扩机器在slaves文件里加上其他节点的主机名然后逐个启动Worker就行。Spark作业提交有local模式和集群模式两种方式。开发调试阶段用local模式最方便spark-submit \ --master local[4] \ --name hypertension-analysis \ analysis.py数据量大之后改成spark-submit \ --master spark://node01:7077 \ --executor-memory 4g \ --total-executor-cores 8 \ analysis.py我实际调试的时候发现executor-memory设太低会导致频繁的GC任务明明能跑但就是慢设太高又容易被YARN或底层物理机限制拒绝。比较稳妥的做法是先用默认参数跑观察Spark UI里各阶段的执行时间再针对性调整。2.4 Spark核心分析逻辑清洗、分级与聚合Spark作业是整个系统里业务逻辑最集中的部分。我用的是PySpark代码结构很清楚读取HDFS上的CSV过滤掉字段缺失和数值异常的数据然后按血压标准分级最后做多维聚合写回MySQL。先看清洗。体检数据脏得很规律常见问题包括收缩压写成“140/90”这样的字符串、舒张压是0、年龄超过120、地区字段为空。这些都要处理掉from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, count, avg, floor spark SparkSession.builder \ .appName(hypertension-analysis) \ .getOrCreate() df spark.read.csv( hdfs://localhost:9000/data/health/*.csv, headerTrue, inferSchemaTrue ) clean_df df.filter( col(systolic).between(30, 250) col(diastolic).between(20, 150) col(age).between(1, 120) col(province).isNotNull() )过滤掉异常值之后做血压分级。这里解释一下为什么不用“家庭自测血压”标准而用“诊室血压”标准——因为这批模拟数据本身标明了测量场景是诊室所以直接按140/90切分更合理。level_df clean_df.withColumn( level, when((col(systolic) 180) | (col(diastolic) 110), 重度) .when((col(systolic) 160) | (col(diastolic) 100), 中度) .when((col(systolic) 140) | (col(diastolic) 90), 轻度) .otherwise(正常) )然后按地区、年龄段、性别做聚合得到患病率。年龄段用floor划分十年一组stat_df level_df.withColumn( age_group, when(col(age) 30, 20-29) .when(col(age) 40, 30-39) .when(col(age) 50, 40-49) .when(col(age) 60, 50-59) .when(col(age) 70, 60-69) .otherwise(70) ).groupBy(province, age_group, gender).agg( count(*).alias(total), count(when(col(level) ! 正常, True)).alias(hypertension_count), avg(systolic).alias(avg_systolic) ) stat_df.withColumn( rate, col(hypertension_count) / col(total) ).write \ .mode(overwrite) \ .jdbc( jdbc:mysql://localhost:3306/hypertension?useUnicodetruecharacterEncodingutf8, region_stat, properties{user: root, password: 123456} )这一段是整个项目的分析核心。水平够的话还可以在这一层加更多指标比如高血压患者中BMI超标占比、吸烟饮酒人群的患病率对比。每个指标都是大屏上的一张图指标越多系统看起来越完整。有一个细节必须提醒Spark写MySQL时如果表里的数据量很大建议先写入临时目录再用DataFrame的overwrite模式覆盖或者用分批写入。否则一次写入几十万行MySQL连接很容易超时。3. Django服务层与可视化大屏实现3.1 Django数据模型设计与ORM查询优化Spark已经把聚合结果写进了MySQLDjango这边就是基于这些表做查询。模型设计上不需要跟Spark输出完全一一对应按大屏展示需求重新组织更好。我建了三个核心模型地区统计表、年龄分布表、危险因素统计表。# models.py from django.db import models class RegionStat(models.Model): province models.CharField(max_length50, verbose_name省份) age_group models.CharField(max_length20, verbose_name年龄段) gender models.CharField(max_length10, verbose_name性别) total models.IntegerField(verbose_name总人数) hypertension_count models.IntegerField(verbose_name高血压人数) rate models.FloatField(verbose_name患病率) avg_systolic models.FloatField(verbose_name平均收缩压) updated_at models.DateTimeField(auto_nowTrue) class Meta: db_table region_statDjango ORM查询这些表很简单但要注意性能。大屏前端每10秒拉一次数据如果每次都把全部记录查出来数据库压力很大。实际开发中我做了两个优化第一用values()只取需要的字段不把整个ORM对象加载到内存data RegionStat.objects.filter(provinceprovince) \ .values(age_group, gender, hypertension_count, total, rate)第二给高频查询字段加索引。Django里用db_indexTrue即可或者直接建Meta.indexes。这个操作在数据量大时差别非常明显不加索引一次查询可能要300ms加了索引能压到几十毫秒。有人会问Django为什么不直接从HDFS读数据因为HDFS不适合高频实时查询而且每次读HDFS都要经过Spark进程延迟极高。正确的做法是先跑Spark生成聚合表然后Django只查MySQL。这个设计是这套系统能流畅运行的关键前提。3.2 后端接口设计与数据推送机制后端接口我用了DRF比手写JsonResponse正规得多。接口设计按大屏的组件拆分区域分布接口、年龄趋势接口、性别对比接口、血压分级占比接口。每个接口都是只读的不涉及写操作。关于Django的查询和删除操作这里单独说一下。ORM里执行查询最常用的是filter()、exclude()、get()删除对象用delete()方法比如# 按条件批量删除 RegionStat.objects.filter(province测试省).delete() # 先查后删返回删除的记录数和明细 obj RegionStat.objects.get(id1024) obj.delete()需要注意get()查不到记录会抛DoesNotExist异常务必用try/except包住或者改用get_object_or_404。批量删除时Django默认会先SELECT出所有符合条件的对象再逐条DELETE数据量特别大的时候建议用_raw_delete()或者直接走原生SQL。数据推送我用了两种方式。第一种是前端轮询最简单兼容性最好setInterval(() { fetch(/api/region/stat/) .then(res res.json()) .then(data updateCharts(data)); }, 10000);第二种是WebSocket推送。后台有新数据写入MySQL后主动通知前端刷新比轮询更及时、更省资源。Django里做WebSocket推荐用Channels核心消费者代码如下# consumers.py import json from channels.generic.websocket import AsyncWebsocketConsumer class DataPushConsumer(AsyncWebsocketConsumer): async def connect(self): await self.accept() await self.channel_layer.group_add(dashboard, self.channel_name) async def disconnect(self, close_code): await self.channel_layer.group_discard(dashboard, self.channel_name) async def data_updated(self, event): await self.send(text_datajson.dumps(event[message]))数据写入方的逻辑是Spark作业完成后调用一个Django接口触发group_send把“数据已更新”的消息广播给所有在线的大屏页面。前端收到消息后重新拉接口图表就能自动刷新。实际测试下来轮询有5-10秒的延迟WebSocket基本是秒级更新体验完全不一样。3.3 可视化大屏的图表布局与核心配置大屏是这套系统的门面布局上我参考了常见的指挥中心风格整体尺寸按1920x1080设计底部留一条滚动数据栏中间核心区域放地图左右两侧放柱状图和饼图。具体到组件顶部大屏标题、当前时间、系统运行状态。中间区域ECharts中国地图用散点图和涟漪效果展示各省高血压患病率。地图数据需要china.js地图包ECharts新版本用registerMap注册GeoJSON。左侧年龄段患病率柱状图、性别占比环形图。右侧血压分级饼图、危险因素Top10横向条形图。底部最近更新的异常记录滚动表格。大屏适配是个容易忽略的坑。直接写死px在分辨率为1366x768的屏幕上会错位所以我用了一个简单的缩放方案把大屏最外层容器固定为1920x1080然后用transform的scale属性按窗口实际大小等比例缩放。意思就是不管屏幕分辨率是多少渲染出来的画面始终保持在1920x1080的视觉比例下。核心ECharts配置片段const mapChart echarts.init(document.getElementById(map)); mapChart.setOption({ tooltip: { trigger: item }, visualMap: { min: 0, max: 30, left: 20, bottom: 20, text: [高患病率, 低患病率], inRange: { color: [#e0f3f8, #abd9e9, #c51b7d] } }, series: [{ type: map, map: china, data: provinceData, label: { show: true, fontSize: 10 } }] });有一点必须提ECharts的visualMap区间如果设置得太死比如固定0到30而实际数据最高只有10那整张地图颜色都会偏淡没有层次。我后来改成动态计算区间用数据最大值的1.2倍作为上限视觉效果立刻好很多。3.4 大屏与后端联调从静态页面到实时数据大屏一开始肯定是静态写死的假数据联调时要逐步替换成后端接口。我习惯分三步调试第一步浏览器直接访问Django接口确认JSON结构第二步在浏览器Console里用fetch测试跨域请求和字段映射第三步把接口数据填入ECharts的series.data里。每替换一个图表就检查一次数据对应关系。跨域是联调阶段最常见的坑。如果大屏页面和后端不在同一台服务器上Django必须配置CORS否则浏览器会直接拦截请求。在Django里用django-cors-headers这个库就能解决但要注意CORS_ALLOW_ALL_ORIGINS True只适合开发环境部署时要改成具体的域名白名单。联调时我踩过最莫名其妙的一个坑是接口返回的数据明明是对的大屏上的地图却显示空白。查了半天才发现是各个图表的数据更新函数全是异步的一个接口报错会导致后续所有.then()里的赋值逻辑都不执行。后来我改成用Promise.all一次拿全部接口数据成功后再统一更新所有图表问题就没了。Promise.all([ fetch(/api/region/stat/).then(r r.json()), fetch(/api/age/trend/).then(r r.json()), fetch(/api/level/ratio/).then(r r.json()) ]).then(([regionData, ageData, levelData]) { updateMap(regionData); updateBar(ageData); updatePie(levelData); });结构化地组织更新逻辑比像挤牙膏一样东一块西一块改代码要省力得多。4. 调试实录、常见问题与性能优化4.1 整条链路的联调流程记录把Hadoop、Spark、Django、大屏四部分单独跑通之后最难啃的是整条链路的联调。我当时的调试流程是这样的先确认HDFS里数据文件能正常读取然后手动跑一遍Spark作业看MySQL里有没有新数据写进来再启动Django接口看数据能不能查到最后启动大屏页面观察图表展示。这套流程里最值得注意的检查点是数据格式的一致性。比如Spark写MySQL的字段类型是VARCHAR还是INTDjango模型里定义的对不对前端parseFloat之后能不能正常使用。类型对不上经常出现“接口返回200但图表不显示”这种看半天发现是NaN的问题。Spark作业跑失败是最常见的联调中断点。我发现线上数据量一大Spark作业就容易因为OOM挂掉。这里提供几个排查方向控制单个executor的内存总量和核数如果数据文件很小就用coalesce把分区数控制在20以内避免小文件导致海量小任务如果某一步join操作特别慢优先考虑广播joinbroadcast join把小表广播到每个executor上。4.2 常见问题速查表以下是我运行这套系统过程中总结的高频问题按频率排序问题现象可能原因解决方法Hadoop DataNode进程起不来反复格式化NameNode导致集群ID不一致删除data目录和tmp目录重新格式化Spark作业报OutOfMemoryErrorexecutor内存不足或分区数过多增大--executor-memory合并小分区Django接口响应很慢没有索引或查了过多字段添加db_index索引用values()只取必要字段大屏地图显示空白地图JSON未加载或异步更新失败确认registerMap已经执行接口用Promise.all集中更新WebSocket页面连不上Channels配置缺少ASGI路由安装daphne配置ASGI_APPLICATIONMySQL写入中文变乱码数据库字符集不是utf8创建库时指定utf8mb4连接串加characterEncodingutf8HDFS上传数据很慢默认副本数过多或者网络带宽限制调低副本数检查节点间网卡速率前端图表出现NaN后端返回字符串类型被直接计算在前端先Number()转换或后端返回float类型这张表基本覆盖了从环境部署到前后联调的大部分坑遇到问题别慌对照现象找原因90%都能在这里找到答案。4.3 性能优化经验与交付建议系统跑通之后如果你想在性能上更进一步有几个性价比很高的优化点第一Spark作业与Web服务彻底分离。跑批和查询都挤在同一台机器上会导致资源争抢最明显的是大屏轮询数据时Spark作业变慢。解决方式是把Spark跑批放到凌晨或者放到另一台独立机器上执行从物理上隔离。第二MySQL查询加Redis缓存。大屏接口是高频只读场景非常适合加缓存。Django里可以用cache_page装饰器或手动在Redis里缓存接口结果5分钟。当数据变化不频繁时缓存能显著降低数据库压力。第三大屏的轮询改成WebSocket 条件触发。如果数据一天只更新一次10秒轮询一次就是极大的浪费。我后来改成每天Spark作业完成后通过WebSocket通知前端刷新更新前后端只发一条消息比轮询省资源响应也更及时。第四前端图表在页面隐藏时暂停刷新。大屏如果长时间挂机可以检测document.hidden状态页面不在视野时暂停轮询回到页面再恢复。这个优化对资源消耗的降低非常明显。关于交付物这些资料在毕业设计里往往就是“源码文档调试过程记录可视化大屏”的打包。文档里我会重点写环境搭建版本兼容性表和部署步骤因为任何人想复现这套系统第一件事就是配环境环境配不上后面全是空谈。调试记录这部分很多同学容易忽略其实非常有价值把遇到的报错和解决过程整理一下文档的专业度直接上一个台阶。5. 扩展思路与个人体会这套系统想往深做还有几个方向可以扩展。接真实数据源的话可以用Flume或FileBeat做实时采集把医院的体检数据自动同步到HDFSSpark改成Structured Streaming做实时流分析加上MLlib还能做高血压风险预测模型把历史病历、家族史、生活方式都喂进去输出个人患病风险概率调度方面接上Airflow让Hadoop跑批、Spark分析、Django刷新整个流程定时自动化。我个人实际调试这套项目的体验是环境搭建花的时间比写业务代码多得多尤其是Hadoop的版本兼容问题Hadoop 3.3和Spark 3.2、Django 4.1之间经常因为JDK版本或者Scala版本不匹配出各种奇怪问题。如果你准备复现这套系统我真心建议在第一台机器上先把全套环境从零搭一遍不要直接复制网上的现成配置因为只有自己踩过一遍环境坑后面分析代码的逻辑才会理解得更透。大屏虽然是最直观的成果展示但真正体现工作量的是那条从HDFS到MySQL的完整数据链路。答辩或者项目汇报时讲清楚“数据从哪来、怎么清洗、怎么算、怎么展示”这条主线比单纯展示一个漂亮的界面要有说服力得多。
返回列表