Spark心脏病数据分析:ETL清洗、分区优化与YARN部署实战解析
2026/9/16 15:03:56 网站建设 项目流程

简介:一份基于Spark的心脏病信息大数据分析毕设源码与配套数据,围绕心脏病相关数据的清洗、处理与特征分析展开,面向计算机相关专业学生及从业者,可作为课程设计、大作业或毕业设计的完整参考。zip压缩包共1010个文件,整体大小8.92MB。文件类型覆盖:Scala/Java等核心分析源码,JS/TS/JSON等前端或项目配置,CSV/Data/XLSX等原始与中间数据,MD/TXT等说明文档,以及JAR/XML/YML等运行环境相关文件,目录结构清晰,便于按模块查阅。目前已有222人学习下载。该项目为个人毕设成果,评审分达97分,代码经过严格调试可运行。内容不仅包含Spark分析流程实现、数据预处理与分区逻辑,还提供了类文件、脚本和说明文档,可帮助读者快速掌握大数据分析项目的工程组织方式,对理解Spark任务编写、数据集划分和结果展示具有直接参考价值。

1. 一个 97 分的 Spark 数据分析毕设,评审真正在看什么

解压这份“基于 Spark 的心脏病信息大数据分析”源码包,你第一眼看到的不是漂亮的可视化,而是一批 Scala object 编译后的 class 文件:ageprocess$.classthalachprocess$.classpartition$.class。这本身就是信号——它拿高分的点不在调参,而在把大数据分析的脏活拆成了可复现的模块。项目要解决的是从原始 CSV 到分析结论之间的清洗、分区、聚合与部署问题,典型的心脏病公开数据集只有几百 KB,但 Spark 的工程框架可以被平移到亿级行为日志。它适合正在选毕业设计方向的学生,也适合刚进数仓团队、想补 Spark 运行时机制的从业者。

2. 从 class 命名反推 ETL:项目结构与数据接入

在 Scala 工程里,object ageprocess编译后会同时生成ageprocess.classageprocess$.class。前者是静态转发入口,后者是模块真正的单例实现。这意味着源码的作者没有把数据处理逻辑堆在一个几百行的 main 里,而是按字段职责拆成独立对象。解压 zip 后看到ageprocessthalachprocesscpprocesspartitionthalach_targetexam这一串名字,基本就能画出这个程序的数据流:读入原始数据 → 字段级清洗 → 分区调整 → 目标构建 → 统计分析。还有一些像hobbysap这类命名,需要打开对应源码看方法签名才能确定业务含义,最好先按这个表定位主要模块。

class 文件推断处理对象典型职责
ageprocess$.classage缺失值处理、年龄分段
thalachprocess$.classthalach最大心率异常值、分箱
cpprocess$.classcp胸痛类型编码
partition$.class分区自定义 Partitioner 或 repartition
thalach_target$.classthalach + target标签联动构造
exam$.class全表数据groupBy 聚合探索

2.1 数据接入:CSV 的 schema 推断开销

读取阶段最常见的写法是spark.read.csvinferSchema,这套源码里的入口大概率也是这样,因为原始数据是带表头的 CSV。代码很容易复现:

val heartDF = spark.read .option("header", "true") // 首行当作列名 .option("inferSchema", "true") // 自动推断字段类型 .csv("hdfs:///user/edu/heart.csv")

header选项控制是否跳过首行,inferSchema决定是否让 Spark 去扫描数据猜测字段类型。这个true是很多 Spark 数据分析案例里用得最顺手、也最容易出问题的地方:inferSchema会触发一次额外的数据扫描,文件小的时候看不出来,一旦数据量到了 GB 级,读取时间会明显变长。

更好的做法是手工声明 schema,尤其当你知道 age、thalach 是整数、target 是 0/1 标签时。代码如下:

import org.apache.spark.sql.types._ val schema = StructType(Array( StructField("age", IntegerType, true), StructField("sex", IntegerType, true), StructField("cp", IntegerType, true), StructField("trestbps", IntegerType, true), StructField("chol", IntegerType, true), StructField("thalach", IntegerType, true), StructField("target", IntegerType, true) )) val heartDF = spark.read .option("header", "true") .schema(schema) // 跳过 inferSchema 的额外扫描 .csv("hdfs:///user/edu/heart.csv")

这里的StructField第三个参数true表示字段允许为空,对应医疗数据里常见的缺失值。如果你发现某些行字段类型被推断成了 string,先别急着 cast,回到原始数据看一眼是不是有空字符串,空字符串转 int 会得到 null,这个细节会直接影响后面的ageprocess判断。

2.2 partition 与 coalesce:分区数不是越大越好

partition$.class的存在说明作者专门处理了分区问题。Spark 里有两个方法容易被混用,repartitioncoalesce。在 ETL 阶段,我一般会这样调整分区:

val balanced = heartDF .repartition(12) // 强制分成 12 个分区 .filter(col("target").isNotNull)

repartition(12)会走一轮全量 shuffle,把数据打散成 12 个分区。这个操作适合在过滤前做,能让后续 stage 的并行度更接近 CPU 核数。但如果只是为了减少输出小文件,应该用coalesce

balanced.coalesce(1) // 从多个分区合并到 1 个,避免 shuffle .write.mode("overwrite") .csv("/tmp/heart_clean")

coalesce(1)只把多个分区合并成少量分区,不触发 shuffle,适合“清理小文件、输出最终结果”的场景。两者差异可以用下面这张表快速判断:

方法是否触发 shuffle适用场景
repartition(n)扩大并行度,重新分布数据
coalesce(n)从多分区往少分区合并

需要注意,coalesce不能指望从 1 个分区扩到 10 个分区,它只擅长缩小。如果看到下游 task 数量明显少于 executor 核数,优先怀疑这里用错了方法。这个分区数的选择也会直接影响后续exam里的groupBy统计速度。

3. 核心清洗链路:age、thalach、cp 的工程实现

拿到ageprocessthalachprocesscpprocess三个 object 后,最重要的其实是方法签名。常见做法是把每个 object 设计成一个接受 DataFrame、返回 DataFrame 的函数式模块,例如def transform(df: DataFrame): DataFrame。这样调用链可以写成cpprocess.transform(thalachprocess.transform(ageprocess.transform(raw))),而这正是后面partitionthalach_target能继续接力下去的前提。

3.1 ageprocess:年龄分段用 when/otherwise,别写 UDF

心脏病数据里 age 几乎不会有负数,但可能存在空值或 0 值。直接把空值删掉会损失样本,一个稳妥的做法是保留成unknown,并在分段时单独统计。下面的逻辑可以用when链实现:

val aged = heartDF .withColumn( "age_group", when(col("age").isNull, "unknown") // 缺失单独成组 .when(col("age") <= 0, "invalid") // 0 值视为异常 .when(col("age") < 30, "20s") .when(col("age") < 40, "30s") .when(col("age") < 50, "40s") .when(col("age") < 60, "50s") .otherwise("60+") )

when链会被 Spark 翻译成类似CASE WHEN的表达式,交给 Catalyst 优化器处理,比注册一个 UDF 快得多。when(...).otherwise(...)的顺序很重要,Spark 从上到下匹配,所以年龄区间必须按从小到大排列,否则 65 岁的人会被提前归到不正确的区间。

3.2 thalachprocess:用 approxQuantile 截断而不是拍脑袋

thalach是最大心率,正常范围大多落在 60~220。用全局maxmin去找范围很容易被个别脏数据带偏,比如把心率记成 300。这里可以用分位数截断:

val Array(low, high) = heartDF .stat .approxQuantile("thalach", Array(0.01, 0.99), 0.0) // 取 1% 和 99% 分位 val thalachDF = aged .withColumn( "thalach_clipped", when(col("thalach") < low, lit(low)) // 低于 1% 分位截断 .when(col("thalach") > high, lit(high)) .otherwise(col("thalach")) )

approxQuantile的第一个参数是列名,第二个参数是要计算的百分位数组,第三个参数是相对误差,0.0 表示精确计算。对于几百 MB 的数据精确计算可以接受;当数据量巨大时,把第三个参数设成 0.05 能显著加快速度,代价是分位点有 5% 的浮动。这样处理后,thalach_clipped比原始字段稳定,后续放进thalach_target作为标签也不容易被异常值干扰。

3.3 cpprocess:先确认类别分布,再决定编码方式

cp(胸痛类型)在 UCI 心脏病数据里是一个 0~3 的类别字段。不同版本的 CSV 编码含义有差别,有的把无症状记为 0,有的记为 4。抄代码之前先跑一条 SQL:

heartDF.groupBy("cp").count().orderBy("cp").show()

如果cpinferSchema读成了字符串,必须先转成整数,否则groupBy之后“0”和“00”会被当成两个类别。转类型很简单:

val cpDF = heartDF .withColumn("cp_int", col("cp").cast("int")) // 字符串转整型

转换后常见的索引结果如下:

cp_int典型含义分析用法
0无症状 / 非心绞痛常作为对照组
1典型心绞痛患病风险高
2非典型心绞痛中等风险
3非心源性疼痛需结合其他指标

如果你的模型是逻辑回归,建议用StringIndexerOneHotEncoder处理;如果只是做 groupBy 统计,整数索引就够了。ageprocessthalachprocess返回的 DataFrame 传给cpprocess时,要注意列名是否重复,withColumn同名会覆盖,而不同名会新增一列,这会影响后面exam里的通配符选择。

4. 从特征到业务结论:exam 与 thalach_target 的验证

清洗完之后,exam承担的是探索性统计。它对应的常见操作是先用createOrReplaceTempView注册临时表,再写 SQL 做分组聚合。与直接在 DataFrame 上调用 API 相比,SQL 在数据集字段多时更容易维护,也方便把某一段逻辑复制到其他分析师工具里。这个案例里最值得关注的业务问题是:患病人群和正常人群在 age、thalach、trestbps 上到底差多少。

4.1 患病率与字段均值:一行 SQL 先看全局

heartDF.createOrReplaceTempView("heartData") val statSQL = """ -- 按 target 分组统计样本量、均值与标准差 SELECT target, COUNT(*) AS cnt, ROUND(AVG(age), 1) AS avg_age, ROUND(AVG(thalach), 1) AS avg_thalach, ROUND(AVG(trestbps), 1) AS avg_trestbps, ROUND(STDDEV(thalach), 1) AS std_thalach FROM heartData GROUP BY target """ spark.sql(statSQL).show()

这段 SQL 的GROUP BY target把样本分成患病(1)和未患病(0)两组,STDDEV计算标准差。只看平均值容易忽略分布差异,比如两组avg_age接近但std_thalach相差很大,说明心率的离散程度才是关键。show()默认只显示前 20 行,这个场景结果只有两行,足够了。

示意输出如下:

targetcntavg_ageavg_thalachstd_thalach
016552.6139.120.8
113848.5149.616.7

在这个公开数据集里,患病组平均年龄更低、平均最大心率更高,这是符合直觉的参考信号。注意这是示意数据,你自己跑出来会有小数点差异,重点看方向。

4.2 thalach_target:把连续特征与标签做成列联表

thalach_target$.class的存在说明源码并不是简单地把 thalach 丢进统计,而是把它和 target 绑定成一个复合标签。一个常见做法是先把thalach >= 150看作“高心率组”,然后和是否患病做交叉表:

val tagged = heartDF .withColumn("thalach_high", when(col("thalach") >= 150, 1).otherwise(0)) .withColumn("is_disease", col("target")) val cross = tagged.stat.crosstab("thalach_high", "is_disease") cross.show()

crosstab的第一个参数是行,第二个参数是列,返回的 DataFrame 第一列名是“第一个参数_第二个参数”。这里得到的是二维频数表,如果thalach_highis_disease不相关,两个组里患病比例应该接近,否则说明该特征值得进入模型。

4.3 用卡方校验给交叉表一个结论

交叉表只能看数字,要判断统计显著性需要卡方检验。Spark MLlib 的ChiSquareTest可以直接在 DataFrame 上跑:

import org.apache.spark.ml.stat.ChiSquareTest val chiDf = tagged .withColumn("feature", col("thalach_high").cast("double")) .select("feature", "is_disease") val chiResult = ChiSquareTest.test(chiDf, "feature", "is_disease") chiResult.show() // 观察 pValue 列

ChiSquareTest.test接收三个参数:样本 DataFrame、特征列名、标签列名。特征必须是数值型,所以这里把thalach_high转成 double。如果输出的pValue小于 0.05,可以认为thalach_high与患病标签有显著关联,这个结论写进毕业设计的“特征有效性验证”章节比单放一张柱状图有说服力得多。

5. spark-submit 部署:内存模型、shuffle 分区与 yarn 客户端误区

拿到源码后很多人会先在本地 IDE 跑通,然后直接丢到集群。此时最容易翻车的是提交参数和内存配置。下面这个提交命令可以作为一个模板,它对应的是 Spark on YARN 的集群模式:

# 向 YARN 提交集群模式作业 spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 4 \ --executor-cores 2 \ --executor-memory 4g \ --driver-memory 2g \ --conf spark.sql.shuffle.partitions=16 \ --class com.example.heart.HeartAnalysis \ heart-analysis.jar \ hdfs:///data/heart.csv hdfs:///output/heart_result

命令最后两个路径会作为args(0)args(1)传给 main 方法,分别对应输入数据和输出目录。关键参数的含义如下:

参数作用常见设定
--num-executors启动的 executor 数量不超过队列限制
--executor-cores单个 executor 占用的 CPU 核数2 或 4
--executor-memory单个 executor 堆内内存根据核数配 2g~8g
--conf spark.sql.shuffle.partitionsshuffle 阶段的默认分区数集群总核数×2 左右

--executor-memory不是越大越好。一个常见的反例是给 8 核的 executor 配 32g 内存,shuffle 时溢出文件会集中在少数几个节点,反而增加 GC 停顿。一般按“每核 2g~4g”预估,比如 4 核对应 8g~16g。

5.1 Spark 内存模型:execution 与 storage 的边界

spark-submit里看到 memory 参数后,还要理解 executor 内部怎么划分内存。Spark 把统一内存池按spark.memory.fraction(默认 0.6)分为执行和存储两部分,剩余 40% 留给用户代码、元数据和字符串池。当groupByjoin这类 shuffle 操作很多时,可以在提交命令里追加:

--conf spark.memory.fraction=0.75 \ --conf spark.memory.storageFraction=0.4

第一个参数调大 execution 可用占比,第二个参数让存储缓存更少占执行内存。但要注意,如果同时用cache()缓存了中间表,存储被挤压会导致重复计算。排查 OOM 时优先看 Spark UI 里每个 Stage 的 “Shuffle Spill (Memory)” 指标,而不是直接加内存。

若收到Container killed by YARN for exceeding memory limits,表示容器总内存超限,此时要增加的是堆外内存占比:

--conf spark.executor.memoryOverhead=1g

memoryOverhead默认是 executor 内存的 10%,在堆内内存不大但序列化、NIO 缓冲较多时非常管用。给 spark 作业加内存的正确顺序是:先看 spill,再调分区数,最后才动 overhead。

5.2 客户端在哪:spark on yarn 只需要一个提交客户端

互联网上经常看到“spark on yarn 提交是不是只需要一个 spark 客户端就行了”这类问题,答案是“是,但客户端要有完整的 Hadoop 配置”。提交时,spark-submit会把 jar 和依赖上传到 HDFS 或临时目录,然后由 ResourceManager 分配 Container,计算并不在提交机器上执行。你需要做的不是给每台机器装 Spark,而是确保提交机器的 spark-env.sh 里有YARN_CONF_DIR,或者yarn-site.xml能被 classpath 找到。

验证提交是否成功,常用命令是:

yarn application -list yarn application -status application_xxx

-status输出里的StateFinalStatus能告诉你作业是 RUNNING、FINISHED 还是 FAILED。如果yarn命令提示找不到 application id,多半是提交到了另一个 YARN 集群或队列,先看--queue参数是否匹配。

6. 把毕设源码改造成 pipeline:自定义分区与 Spark UI 验证

这套源码的 object 拆分已经接近 pipeline 的雏形。真正落地时值得做两件事:把多个 process 串成一个 main,并自定义分区让数据分布更可控。

6.1 串成函数式调用链

每个 process object 都设计成transform(df: DataFrame): DataFrame后,主程序可以很干净:

object HeartAnalysis { def main(args: Array[String]): Unit = { val spark = SparkSession.builder .appName("HeartAnalysis") .getOrCreate() val input = args(0) val output = args(1) val raw = spark.read.option("header", "true").csv(input) val result = cpprocess.transform( thalachprocess.transform( ageprocess.transform(raw))) result .coalesce(1) .write.mode("overwrite") .parquet(output) } }

这里coalesce(1)只适合最终结果很小的情况。如果要把中间结果继续给下游用,建议write.parquet时保留原生分区数,并按照常用过滤字段做分区:

result.write .partitionBy("target") // 按标签值分目录存储 .mode("overwrite") .parquet(output)

partitionBy("target")会为 0/1 各写一个目录,后续查询只扫一个目录而不是全表。注意这个参数和repartition不同,它改的是存储结构,不直接提升单次计算并行度。

6.2 自定义分区的最终验证

如果项目里的partition$.class用的是 RDD API,那大概率是继承了Partitioner。一个按年龄段分桶的实现长这样:

class AgeBucketPartitioner(numParts: Int) extends Partitioner { override def numPartitions: Int = numParts override def getPartition(key: Any): Int = key match { case age: Int => math.min(age / 10, numParts - 1) // 每 10 岁一个桶 case _ => math.abs(key.hashCode()) % numParts } }

numPartitions定义桶数,getPartition返回 key 所在桶。使用时先取 RDD,再调用partitionBy

val rdd = spark.sparkContext.parallelize(heartDF.rdd.map(r => (r.getAs[Int]("age"), r))) val bucketed = rdd.partitionBy(new AgeBucketPartitioner(6))

partitionBy会发生一次 shuffle,但后续同一区的数据在同一个 executor 上,适合“按年龄分文件输出”这类需求。验证分区是否均衡,打开 Spark UI 的 Stages 页面,看每个 Task 的 Shuffle Write Size 和 Duration 是否接近;如果某个 Task 的数据量是其他的两倍以上,说明分区键选择不当,回到getPartition里加一个年龄分布统计,再调整桶边界。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询