【大数据】肝硬化患者生存预测数据可视化分析系统 Hadoop+Spark技术 计算机毕业设计项目 Anaconda环境配置 附源码+文档+讲解
一、个人简介
💖💖作者:计算机编程果茶熊
💙💙个人简介:曾长期从事计算机专业培训教学,担任过编程老师,同时本人也热爱上课教学,擅长Java、微信小程序、Python、Golang、安卓Android等多个IT方向。会做一些项目定制化开发、代码讲解、答辩教学、文档编写、也懂一些降重方面的技巧。平常喜欢分享一些自己开发中遇到的问题的解决办法,也喜欢交流技术,大家有技术代码这一块的问题可以问我!
💛💛想说的话:感谢大家的关注与支持!
💜💜
网站实战项目
安卓/小程序实战项目
大数据实战项目
计算机毕业设计选题
💕💕文末获取源码联系计算机编程果茶熊
二、系统介绍
大数据框架:Hadoop+Spark(本次没用Hive,支持定制)
开发语言:Python
后端框架:Django
前端:Vue
详细技术点:Hadoop、HDFS、Spark、Spark SQL、Pandas、NumPy
数据库:MySQL
基于大数据的肝硬化患者生存预测数据可视化分析系统采用Hadoop与Spark作为底层大数据处理框架,通过HDFS实现海量医疗数据的分布式存储,利用Spark SQL完成对患者生化指标、临床分期、人口特征等多维度数据的快速查询与计算。系统后端基于Django框架搭建,前端采用Vue技术栈开发交互界面,结合Pandas和NumPy进行数据预处理与统计分析,MySQL数据库负责存储用户信息及系统配置数据。平台提供了患者数据管理、生化指标分析、临床分期分析、人口特征分析等核心功能模块,通过多维交叉特征分析挖掘患者病情发展的潜在规律,借助生存曲线分析展示不同患者群体的预后情况,治疗方案预后分析模块能够对比各类治疗手段的效果差异。数据大屏以可视化图表的形式呈现关键指标的统计结果,管理员可以通过用户管理与系统公告管理模块维护平台的日常运行,整个系统在处理大规模肝硬化患者数据时展现出较好的性能表现,为医疗机构开展患者预后评估提供了技术支撑。
三、视频解说
四、部分功能展示







五、部分代码展示
from pyspark.sql import SparkSession
from pyspark.sql.functions import col,count,avg,sum,when,datediff,current_date,year,month,desc,asc,round
from pyspark.sql.types import StructType,StructField,StringType,IntegerType,FloatType,DateType
from django.http import JsonResponse
from django.views.decorators.http import require_http_methods
import pandas as pd
import numpy as np
from datetime import datetime
import json
spark=SparkSession.builder.appName("LiverCirrhosisAnalysis").config("spark.sql.warehouse.dir","/user/hive/warehouse").config("spark.executor.memory","4g").config("spark.driver.memory","2g").getOrCreate()
@require_http_methods(["POST"])
def analyze_biochemical_indicators(request):
try:
request_data=json.loads(request.body)
patient_ids=request_data.get('patient_ids',[])
indicator_types=request_data.get('indicator_types',['ALT','AST','TBIL','ALB','PT'])
time_range_start=request_data.get('start_date')
time_range_end=request_data.get('end_date')
hdfs_path="hdfs://namenode:9000/medical_data/biochemical_records"
biochemical_df=spark.read.parquet(hdfs_path)
filtered_df=biochemical_df.filter((col("test_date")>=time_range_start)&(col("test_date")<=time_range_end))
if patient_ids:
filtered_df=filtered_df.filter(col("patient_id").isin(patient_ids))
indicator_stats={}
for indicator in indicator_types:
indicator_col=indicator.lower()+"_value"
stats_df=filtered_df.select(col("patient_id"),col(indicator_col),col("test_date")).groupBy("patient_id").agg(avg(indicator_col).alias("avg_value"),count(indicator_col).alias("test_count"),round(avg(indicator_col),2).alias("mean_value"))
abnormal_threshold={"ALT":40,"AST":40,"TBIL":21,"ALB":35,"PT":13}
threshold_value=abnormal_threshold.get(indicator,0)
abnormal_df=filtered_df.filter(col(indicator_col)>threshold_value).groupBy("patient_id").agg(count("*").alias("abnormal_count"))
merged_df=stats_df.join(abnormal_df,on="patient_id",how="left").fillna(0)
trend_df=filtered_df.select("patient_id",indicator_col,"test_date").orderBy("patient_id","test_date")
window_spec=Window.partitionBy("patient_id").orderBy("test_date")
trend_df=trend_df.withColumn("prev_value",lag(indicator_col,1).over(window_spec))
trend_df=trend_df.withColumn("value_change",when(col("prev_value").isNotNull(),col(indicator_col)-col("prev_value")).otherwise(0))
trend_summary=trend_df.groupBy("patient_id").agg(avg("value_change").alias("avg_change"),sum(when(col("value_change")>0,1).otherwise(0)).alias("increase_times"))
final_result=merged_df.join(trend_summary,on="patient_id",how="left")
result_pandas=final_result.toPandas()
indicator_stats[indicator]=result_pandas.to_dict(orient='records')
correlation_matrix={}
if len(indicator_types)>=2:
correlation_cols=[ind.lower()+"_value" for ind in indicator_types]
correlation_df=filtered_df.select(correlation_cols).toPandas()
corr_result=correlation_df.corr()
correlation_matrix=corr_result.to_dict()
response_data={"status":"success","indicator_statistics":indicator_stats,"correlation_matrix":correlation_matrix,"total_patients":filtered_df.select("patient_id").distinct().count(),"analysis_period":{"start":time_range_start,"end":time_range_end}}
return JsonResponse(response_data,safe=False)
except Exception as e:
return JsonResponse({"status":"error","message":str(e)},status=500)
@require_http_methods(["POST"])
def analyze_survival_curve(request):
try:
request_data=json.loads(request.body)
grouping_factor=request_data.get('grouping_factor','child_pugh_grade')
follow_up_months=request_data.get('follow_up_months',60)
include_censored=request_data.get('include_censored',True)
hdfs_patient_path="hdfs://namenode:9000/medical_data/patient_records"
hdfs_followup_path="hdfs://namenode:9000/medical_data/followup_records"
patient_df=spark.read.parquet(hdfs_patient_path)
followup_df=spark.read.parquet(hdfs_followup_path)
merged_df=patient_df.join(followup_df,on="patient_id",how="inner")
merged_df=merged_df.withColumn("survival_months",datediff(col("last_followup_date"),col("diagnosis_date"))/30)
merged_df=merged_df.filter(col("survival_months")<=follow_up_months)
if not include_censored:
merged_df=merged_df.filter(col("survival_status")=="deceased")
grouped_data=merged_df.groupBy(grouping_factor,"survival_months").agg(count("patient_id").alias("event_count"))
total_patients_per_group=merged_df.groupBy(grouping_factor).agg(count("patient_id").alias("total_count"))
survival_data=grouped_data.join(total_patients_per_group,on=grouping_factor,how="inner")
window_spec=Window.partitionBy(grouping_factor).orderBy("survival_months").rowsBetween(Window.unboundedPreceding,Window.currentRow)
survival_data=survival_data.withColumn("cumulative_events",sum("event_count").over(window_spec))
survival_data=survival_data.withColumn("survival_rate",round((col("total_count")-col("cumulative_events"))/col("total_count"),4))
survival_data=survival_data.withColumn("mortality_rate",round(col("cumulative_events")/col("total_count"),4))
survival_curves={}
group_values=survival_data.select(grouping_factor).distinct().collect()
for group_row in group_values:
group_value=group_row[grouping_factor]
group_df=survival_data.filter(col(grouping_factor)==group_value).orderBy("survival_months")
group_pandas=group_df.select("survival_months","survival_rate","mortality_rate","event_count").toPandas()
survival_curves[group_value]=group_pandas.to_dict(orient='records')
median_survival_times={}
for group_value in survival_curves.keys():
curve_data=survival_curves[group_value]
median_time=None
for point in curve_data:
if point['survival_rate']<=0.5:
median_time=point['survival_months']
break
median_survival_times[group_value]=median_time if median_time else "Not reached"
hazard_ratios={}
if grouping_factor=="child_pugh_grade":
grade_a_df=merged_df.filter(col(grouping_factor)=="A")
grade_b_df=merged_df.filter(col(grouping_factor)=="B")
grade_c_df=merged_df.filter(col(grouping_factor)=="C")
a_death_rate=grade_a_df.filter(col("survival_status")=="deceased").count()/grade_a_df.count()
b_death_rate=grade_b_df.filter(col("survival_status")=="deceased").count()/grade_b_df.count()
c_death_rate=grade_c_df.filter(col("survival_status")=="deceased").count()/grade_c_df.count()
hazard_ratios["B_vs_A"]=round(b_death_rate/a_death_rate,3) if a_death_rate>0 else None
hazard_ratios["C_vs_A"]=round(c_death_rate/a_death_rate,3) if a_death_rate>0 else None
response_data={"status":"success","survival_curves":survival_curves,"median_survival_times":median_survival_times,"hazard_ratios":hazard_ratios,"grouping_factor":grouping_factor,"follow_up_months":follow_up_months}
return JsonResponse(response_data,safe=False)
except Exception as e:
return JsonResponse({"status":"error","message":str(e)},status=500)
@require_http_methods(["POST"])
def analyze_multidimensional_features(request):
try:
request_data=json.loads(request.body)
dimension_list=request_data.get('dimensions',['age_group','gender','child_pugh_grade'])
target_indicator=request_data.get('target_indicator','survival_months')
aggregation_method=request_data.get('aggregation','avg')
hdfs_integrated_path="hdfs://namenode:9000/medical_data/integrated_patient_data"
integrated_df=spark.read.parquet(hdfs_integrated_path)
integrated_df=integrated_df.withColumn("age_group",when(col("age")<40,"<40").when((col("age")>=40)&(col("age")<60),"40-60").otherwise(">=60"))
integrated_df=integrated_df.withColumn("bilirubin_level",when(col("total_bilirubin")<34,"Normal").when((col("total_bilirubin")>=34)&(col("total_bilirubin")<51),"Mild").otherwise("Severe"))
cross_analysis_results={}
for i in range(len(dimension_list)):
for j in range(i+1,len(dimension_list)):
dim1=dimension_list[i]
dim2=dimension_list[j]
cross_key=f"{dim1}_vs_{dim2}"
if aggregation_method=="avg":
agg_func=avg(target_indicator)
elif aggregation_method=="sum":
agg_func=sum(target_indicator)
elif aggregation_method=="count":
agg_func=count(target_indicator)
else:
agg_func=avg(target_indicator)
cross_df=integrated_df.groupBy(dim1,dim2).agg(agg_func.alias("aggregated_value"),count("patient_id").alias("patient_count"))
cross_df=cross_df.withColumn("percentage",round(col("patient_count")/sum("patient_count").over(Window.partitionBy()),4))
cross_pandas=cross_df.orderBy(desc("aggregated_value")).toPandas()
cross_analysis_results[cross_key]=cross_pandas.to_dict(orient='records')
three_dim_analysis={}
if len(dimension_list)>=3:
dim1,dim2,dim3=dimension_list[0],dimension_list[1],dimension_list[2]
three_dim_key=f"{dim1}_{dim2}_{dim3}"
three_dim_df=integrated_df.groupBy(dim1,dim2,dim3).agg(avg(target_indicator).alias("avg_value"),count("patient_id").alias("count"),sum(when(col("survival_status")=="deceased",1).otherwise(0)).alias("death_count"))
three_dim_df=three_dim_df.withColumn("mortality_rate",round(col("death_count")/col("count"),4))
three_dim_pandas=three_dim_df.orderBy(desc("mortality_rate")).toPandas()
three_dim_analysis[three_dim_key]=three_dim_pandas.to_dict(orient='records')
feature_importance={}
for dimension in dimension_list:
dim_variance=integrated_df.groupBy(dimension).agg(avg(target_indicator).alias("group_avg"))
overall_avg=integrated_df.agg(avg(target_indicator)).collect()[0][0]
variance_sum=dim_variance.withColumn("variance",(col("group_avg")-overall_avg)**2).agg(sum("variance")).collect()[0][0]
feature_importance[dimension]=round(variance_sum,4)
sorted_importance=dict(sorted(feature_importance.items(),key=lambda x:x[1],reverse=True))
heatmap_data={}
if len(dimension_list)==2:
dim1,dim2=dimension_list[0],dimension_list[1]
heatmap_df=integrated_df.groupBy(dim1,dim2).agg(avg(target_indicator).alias("value"))
heatmap_pandas=heatmap_df.toPandas()
pivot_table=heatmap_pandas.pivot(index=dim1,columns=dim2,values='value')
heatmap_data=pivot_table.to_dict()
response_data={"status":"success","cross_analysis":cross_analysis_results,"three_dimensional_analysis":three_dim_analysis,"feature_importance":sorted_importance,"heatmap_data":heatmap_data,"dimensions_analyzed":dimension_list,"target_indicator":target_indicator}
return JsonResponse(response_data,safe=False)
except Exception as e:
return JsonResponse({"status":"error","message":str(e)},status=500)
六、部分文档展示

七、END
💕💕文末获取源码联系计算机编程果茶熊
更多推荐



所有评论(0)