You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.31 12:55:34