Spark与JVM对象内存开销疑问及Tungsten优化探究
Spark内存管理与Tungsten优化:对象内存开销解析
问题背景
我正在研究Apache Spark的内存管理,试图理解对象存储的内存开销问题——毕竟Spark运行在JVM之上。同事认为Spark基于JVM会存在大量内存开销,举例说Spark里每个单字符字符串会因JVM开销占用40字节以上内存。但我觉得借助Tungsten这类优化技术,Spark用特殊格式存储数据,其中的字符串和JVM原生字符串并不等价。
为验证这个问题,我尝试测算DataFrame的大小,但结果超出预期:
val df = spark.range(100000000).selectExpr("cast(id as string) as x").repartition(1) SizeEstimator.estimate(df) // 结果为4mb val df2 = spark.range(100000000).selectExpr("concat(cast(id as string),'z') as x").repartition(1) SizeEstimator.estimate(df2) // 结果同样为4mb
两个DataFrame理论上内存占用应该不同(一个无额外字符,一个多了一个字符),但估算结果完全一致。
问题解答
1. Spark场景下JVM对象内存开销的认知准确性与优化程度
- 同事的说法仅适用于未优化的原生RDD:比如
JavaRDD[String]里的每个JVM String对象,确实会有高额额外开销——64位JVM默认压缩指针下,对象头占12字节,char数组的对象头又占12字节,单字符的char数组占2字节,再加上内存对齐,单字符String确实可能占用40字节左右。 - 但对于DataFrame/Dataset,Spark做了根本性优化,完全规避了原生JVM对象的高额开销:Spark会将数据存储在基于Tungsten的列存格式中,而非原生JVM对象集合,所以不存在传统JVM对象的内存冗余问题。
2. Tungsten等优化技术对DataFrame/RDD内存占用的影响
Tungsten是Spark的核心优化引擎,主要通过以下方式降低内存占用:
- 二进制列存格式:数据以紧凑的二进制形式按列存储,而非JVM对象。比如字符串会被编码成UTF-8字节数组直接存入内存块,没有JVM对象头、数组头这类额外开销。像你测试的场景,数字字符串和加了'z'的字符串,会以实际字节长度存储,而非JVM String的臃肿结构。
- 直接内存操作:借助
sun.misc.Unsafe直接操作内存,绕过JVM的自动内存管理,同时保证内存布局的连续性,减少内存碎片化。 - 高效编码策略:针对不同数据类型使用专属编码,比如对重复字符串用字典编码,对数字字符串做压缩处理,进一步降低内存占用。
关于SizeEstimator结果异常的说明
你用SizeEstimator.estimate(df)得到的4MB结果,并非DataFrame的实际内存占用,原因如下:
SizeEstimator本质是估算原生RDD的内存占用,但DataFrame的底层是逻辑执行计划,并非直接的JVM对象集合,它无法准确识别Tungsten列存格式的实际内存大小。- 要准确测算DataFrame内存占用,可以:
- 触发Action(比如执行
df.cache()后,查看Spark UI的Storage页面),这里会显示实际缓存到内存的大小。 - 将DataFrame转为RDD后再用
SizeEstimator,但结果仍有偏差,因为RDD的存储形式和DataFrame的列存结构不同。
- 触发Action(比如执行
内容的提问来源于stack exchange,提问作者HideDJeker
相关产品推荐
相关产品推荐

