前言

💖💖作者:计算机程序员小杨
💙💙个人简介:我是一名计算机相关专业的从业者,擅长Java、微信小程序、Python、Golang、安卓Android等多个IT方向。会做一些项目定制化开发、代码讲解、答辩教学、文档编写、也懂一些降重方面的技巧。热爱技术,喜欢钻研新工具和框架,也乐于通过代码解决实际问题,大家有技术代码这一块的问题可以问我!
💛💛想说的话:感谢大家的关注与支持!
💕💕文末获取源码联系 计算机程序员小杨
💜💜
网站实战项目
安卓/小程序实战项目
大数据实战项目
深度学习实战项目
计算机毕业设计选题
💜💜

一.开发工具简介

大数据框架:Hadoop+Spark(本次没用Hive,支持定制)
开发语言:Python
后端框架:Django
前端:Vue
详细技术点:Hadoop、HDFS、Spark、Spark SQL、Pandas、NumPy
数据库:MySQL

二.系统内容简介

基于大数据的高血压风险数据可视化分析系统通过整合Hadoop与Spark分布式计算框架,实现了对海量高血压相关数据的高效处理与深度挖掘。系统采用Django作为后端服务框架,结合Vue构建交互友好的前端界面,在HDFS分布式文件系统中存储原始数据,利用Spark SQL完成数据清洗、转换等预处理工作,通过Pandas与NumPy进行统计分析与特征提取。系统提供了用户管理、高血压风险数据管理等基础功能模块,在数据分析层面设计了基础用户画像分析、生活习惯风险分析、多维综合探索分析以及核心生理指标分析等多个维度的分析模块,通过系统大屏将分析结果以图表形式直观呈现,帮助医疗工作者或研究人员从不同角度理解高血压风险因素的分布特征与关联规律,为高血压的预防与干预提供数据支撑,同时也为大数据技术在医疗健康领域的应用探索了一条可行路径。

三.系统功能演示

基于大数据的高血压风险数据可视化分析系统–演示视频

四.系统界面展示

在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述
在这里插入图片描述

五.系统源码展示

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, avg, count, when, sum as spark_sum, round as spark_round
from django.http import JsonResponse
from django.views import View
import pandas as pd
import numpy as np
import json

spark = SparkSession.builder.appName("HypertensionRiskAnalysis").config("spark.sql.warehouse.dir", "/user/hive/warehouse").config("spark.executor.memory", "2g").config("spark.driver.memory", "1g").getOrCreate()

class UserPortraitAnalysisView(View):
    def get(self, request):
        try:
            hdfs_path = "hdfs://localhost:9000/hypertension_data/user_health_records.csv"
            df = spark.read.csv(hdfs_path, header=True, inferSchema=True)
            df = df.filter(col("age").isNotNull() & col("gender").isNotNull() & col("systolic_pressure").isNotNull())
            age_groups = df.withColumn("age_group", when(col("age") < 30, "30岁以下").when((col("age") >= 30) & (col("age") < 40), "30-40岁").when((col("age") >= 40) & (col("age") < 50), "40-50岁").when((col("age") >= 50) & (col("age") < 60), "50-60岁").otherwise("60岁以上"))
            age_distribution = age_groups.groupBy("age_group").agg(count("*").alias("count"), avg("systolic_pressure").alias("avg_systolic"), avg("diastolic_pressure").alias("avg_diastolic"))
            age_distribution = age_distribution.withColumn("avg_systolic", spark_round(col("avg_systolic"), 2)).withColumn("avg_diastolic", spark_round(col("avg_diastolic"), 2))
            age_result = age_distribution.orderBy("age_group").collect()
            age_data = [{"age_group": row["age_group"], "count": row["count"], "avg_systolic": float(row["avg_systolic"]), "avg_diastolic": float(row["avg_diastolic"])} for row in age_result]
            gender_distribution = df.groupBy("gender").agg(count("*").alias("count"), avg("systolic_pressure").alias("avg_systolic"))
            gender_distribution = gender_distribution.withColumn("avg_systolic", spark_round(col("avg_systolic"), 2))
            gender_result = gender_distribution.collect()
            gender_data = [{"gender": "男性" if row["gender"] == "M" else "女性", "count": row["count"], "avg_systolic": float(row["avg_systolic"])} for row in gender_result]
            high_risk_count = df.filter((col("systolic_pressure") >= 140) | (col("diastolic_pressure") >= 90)).count()
            total_count = df.count()
            risk_ratio = round((high_risk_count / total_count) * 100, 2) if total_count > 0 else 0
            response_data = {"age_distribution": age_data, "gender_distribution": gender_data, "high_risk_ratio": risk_ratio, "total_users": total_count}
            return JsonResponse({"code": 200, "message": "用户画像分析完成", "data": response_data})
        except Exception as e:
            return JsonResponse({"code": 500, "message": f"分析过程出现错误: {str(e)}"})

class LifestyleRiskAnalysisView(View):
    def get(self, request):
        try:
            hdfs_path = "hdfs://localhost:9000/hypertension_data/lifestyle_records.csv"
            df = spark.read.csv(hdfs_path, header=True, inferSchema=True)
            df = df.filter(col("smoking_status").isNotNull() & col("drinking_frequency").isNotNull() & col("exercise_hours").isNotNull())
            smoking_analysis = df.groupBy("smoking_status").agg(count("*").alias("count"), avg("systolic_pressure").alias("avg_pressure"), spark_sum(when(col("systolic_pressure") >= 140, 1).otherwise(0)).alias("high_risk_count"))
            smoking_analysis = smoking_analysis.withColumn("avg_pressure", spark_round(col("avg_pressure"), 2))
            smoking_analysis = smoking_analysis.withColumn("risk_percentage", spark_round((col("high_risk_count") / col("count")) * 100, 2))
            smoking_result = smoking_analysis.collect()
            smoking_data = [{"status": "吸烟" if row["smoking_status"] == 1 else "不吸烟", "count": row["count"], "avg_pressure": float(row["avg_pressure"]), "risk_percentage": float(row["risk_percentage"])} for row in smoking_result]
            drinking_analysis = df.groupBy("drinking_frequency").agg(count("*").alias("count"), avg("systolic_pressure").alias("avg_pressure"))
            drinking_analysis = drinking_analysis.withColumn("avg_pressure", spark_round(col("avg_pressure"), 2))
            drinking_result = drinking_analysis.orderBy("drinking_frequency").collect()
            frequency_map = {0: "不饮酒", 1: "偶尔", 2: "经常", 3: "频繁"}
            drinking_data = [{"frequency": frequency_map.get(row["drinking_frequency"], "未知"), "count": row["count"], "avg_pressure": float(row["avg_pressure"])} for row in drinking_result]
            exercise_bins = [0, 2, 5, 10, 24]
            exercise_labels = ["2小时以下", "2-5小时", "5-10小时", "10小时以上"]
            pandas_df = df.select("exercise_hours", "systolic_pressure").toPandas()
            pandas_df["exercise_group"] = pd.cut(pandas_df["exercise_hours"], bins=exercise_bins, labels=exercise_labels, include_lowest=True)
            exercise_stats = pandas_df.groupby("exercise_group", observed=True).agg(count=("exercise_hours", "count"), avg_pressure=("systolic_pressure", "mean")).reset_index()
            exercise_data = [{"group": row["exercise_group"], "count": int(row["count"]), "avg_pressure": round(row["avg_pressure"], 2)} for _, row in exercise_stats.iterrows()]
            response_data = {"smoking_analysis": smoking_data, "drinking_analysis": drinking_data, "exercise_analysis": exercise_data}
            return JsonResponse({"code": 200, "message": "生活习惯风险分析完成", "data": response_data})
        except Exception as e:
            return JsonResponse({"code": 500, "message": f"分析过程出现错误: {str(e)}"})

class PhysiologicalIndicatorAnalysisView(View):
    def post(self, request):
        try:
            params = json.loads(request.body)
            indicator_type = params.get("indicator_type", "blood_pressure")
            time_range = params.get("time_range", 30)
            hdfs_path = "hdfs://localhost:9000/hypertension_data/physiological_indicators.csv"
            df = spark.read.csv(hdfs_path, header=True, inferSchema=True)
            df = df.filter(col("record_date").isNotNull())
            from pyspark.sql.functions import datediff, current_date
            df = df.filter(datediff(current_date(), col("record_date")) <= time_range)
            if indicator_type == "blood_pressure":
                pressure_ranges = df.withColumn("pressure_level", when(col("systolic_pressure") < 120, "正常").when((col("systolic_pressure") >= 120) & (col("systolic_pressure") < 140), "偏高").otherwise("高血压"))
                distribution = pressure_ranges.groupBy("pressure_level").agg(count("*").alias("count"), avg("systolic_pressure").alias("avg_systolic"), avg("diastolic_pressure").alias("avg_diastolic"))
                distribution = distribution.withColumn("avg_systolic", spark_round(col("avg_systolic"), 2)).withColumn("avg_diastolic", spark_round(col("avg_diastolic"), 2))
                result = distribution.collect()
                analysis_data = [{"level": row["pressure_level"], "count": row["count"], "avg_systolic": float(row["avg_systolic"]), "avg_diastolic": float(row["avg_diastolic"])} for row in result]
                correlation_df = df.select("systolic_pressure", "heart_rate", "bmi").toPandas()
                correlation_matrix = correlation_df.corr()
                correlation_data = {"pressure_heart_rate": round(correlation_matrix.loc["systolic_pressure", "heart_rate"], 3), "pressure_bmi": round(correlation_matrix.loc["systolic_pressure", "bmi"], 3)}
            elif indicator_type == "heart_rate":
                heart_rate_ranges = df.withColumn("heart_rate_level", when(col("heart_rate") < 60, "偏低").when((col("heart_rate") >= 60) & (col("heart_rate") <= 100), "正常").otherwise("偏高"))
                distribution = heart_rate_ranges.groupBy("heart_rate_level").agg(count("*").alias("count"), avg("heart_rate").alias("avg_heart_rate"))
                distribution = distribution.withColumn("avg_heart_rate", spark_round(col("avg_heart_rate"), 2))
                result = distribution.collect()
                analysis_data = [{"level": row["heart_rate_level"], "count": row["count"], "avg_heart_rate": float(row["avg_heart_rate"])} for row in result]
                correlation_data = {}
            else:
                bmi_ranges = df.withColumn("bmi_level", when(col("bmi") < 18.5, "偏瘦").when((col("bmi") >= 18.5) & (col("bmi") < 24), "正常").when((col("bmi") >= 24) & (col("bmi") < 28), "偏胖").otherwise("肥胖"))
                distribution = bmi_ranges.groupBy("bmi_level").agg(count("*").alias("count"), avg("bmi").alias("avg_bmi"), avg("systolic_pressure").alias("avg_pressure"))
                distribution = distribution.withColumn("avg_bmi", spark_round(col("avg_bmi"), 2)).withColumn("avg_pressure", spark_round(col("avg_pressure"), 2))
                result = distribution.collect()
                analysis_data = [{"level": row["bmi_level"], "count": row["count"], "avg_bmi": float(row["avg_bmi"]), "avg_pressure": float(row["avg_pressure"])} for row in result]
                correlation_data = {}
            response_data = {"indicator_type": indicator_type, "time_range": time_range, "distribution": analysis_data, "correlation": correlation_data}
            return JsonResponse({"code": 200, "message": "生理指标分析完成", "data": response_data})
        except Exception as e:
            return JsonResponse({"code": 500, "message": f"分析过程出现错误: {str(e)}"})

六.系统文档展示

在这里插入图片描述

结束

💕💕文末获取源码联系 计算机程序员小杨

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐