Spark 3.3.1中MEMORY_AND_DISK持久化RDD丢失问题咨询及示例请求
问题原因分析
核心矛盾:内存资源限制与序列化存储的额外开销
你的集群配置为1个Executor(3核、4GB内存),结合数据特性,问题根源如下:
- 内存可用空间的实际限制:Executor的4GB内存并非全部用于缓存,默认会预留10%作为
spark.executor.memoryOverhead(约409MB),实际可用缓存内存约为3.59GB。 MEMORY_AND_DISK的内存额外消耗:该存储级别会先尝试将序列化后的分区存入内存,但序列化过程中需维护RDD元数据、分区索引等结构,实际内存占用会高于单纯的序列化数据大小(1965.2MB)。当多分区并发写入内存时,易瞬间突破内存阈值,触发OOM或YARN容器被Kill,对应YarnAllocator: Container from a bad node告警。- 单Executor的容错缺失:仅1个Executor且默认副本数为1,一旦节点异常,所有缓存块无副本可恢复,就会出现
No more replicas available for rdd_13_3告警。 - 与
MEMORY_AND_DISK_DESER的差异:反序列化存储的3.6GB数据刚好接近可用内存上限,且无需额外的序列化临时内存开销,Spark可更高效管理反序列化对象,因此能稳定运行。
可复现示例代码
以下代码生成与源数据规模匹配的随机CSV,按你的集群配置读取后,使用MEMORY_AND_DISK(DISK_AND_MEMORY为其别名)持久化,会同时用到内存与磁盘:
步骤1:生成测试CSV
import pyspark from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType import random # 临时创建SparkSession用于生成数据 spark = SparkSession.builder.appName("GenerateTestData").getOrCreate() # 模拟CDC COVID数据的Schema schema = StructType([ StructField("case_month", StringType(), True), StructField("res_state", StringType(), True), StructField("age_group", StringType(), True), StructField("sex", StringType(), True), StructField("hosp_yn", StringType(), True), StructField("death_yn", StringType(), True), StructField("case_positive_specimen_interval", IntegerType(), True) ]) # 生成单条随机数据 def gen_random_row(): months = ["2020-01", "2020-06", "2021-01", "2021-06"] states = ["CA", "TX", "FL", "NY", "IL"] age_groups = ["0-9", "20-29", "40-49", "60-69", "80+"] yes_no_unknown = ["Yes", "No", "Unknown"] return ( random.choice(months), random.choice(states), random.choice(age_groups), random.choice(["Male", "Female", "Unknown"]), random.choice(yes_no_unknown), random.choice(yes_no_unknown), random.randint(0, 30) ) # 生成约12GB的数据集(单条数据约150字节,8亿行) rdd = spark.sparkContext.parallelize(range(80000000), numSlices=100).map(lambda x: gen_random_row()) df = spark.createDataFrame(rdd, schema) # 写入CSV(替换为你的存储路径) df.write.mode("overwrite").option("header", "true").csv("/path/to/test_covid_data") spark.stop()
步骤2:按指定配置读取并持久化
import pyspark from pyspark.sql import SparkSession from pyspark import StorageLevel # 按你的集群配置初始化SparkSession spark = SparkSession.builder.\ appName("TestMemoryDiskPersistence").\ config("spark.dynamicAllocation.enabled", False).\ config("spark.executor.cores", "3").\ config("spark.executor.instances", "1").\ config("spark.executor.memory", "4G").\ config("spark.sql.adaptive.enabled", True).\ config("spark.storage.blockManagerSlaveTimeoutMs", "300000").\ # 延长超时避免误判节点异常 getOrCreate() # 读取测试CSV df = spark.read.option("header", "true").csv("/path/to/test_covid_data") # 使用MEMORY_AND_DISK持久化 df.persist(StorageLevel.MEMORY_AND_DISK) # 触发持久化动作 df.count() # 查看缓存状态 print("缓存级别:", df.storageLevel) storage_status = spark.sparkContext.getStorageStatus()[0] print(f"内存缓存大小:{storage_status.memSize / 1024 / 1024:.2f} MB") print(f"磁盘缓存大小:{storage_status.diskSize / 1024 / 1024:.2f} MB") spark.stop()
优化建议
- 增加
spark.executor.memoryOverhead至1GB,避免内存溢出:config("spark.executor.memoryOverhead", "1G") - 若集群资源允许,增加Executor数量,提升容错性与缓存能力
- 多Executor场景下,将
spark.storage.replicationFactor设为2,防止单节点异常丢失数据
内容的提问来源于stack exchange,提问作者figs_and_nuts
相关产品推荐
相关产品推荐

