在AWS EMR不同步骤中能否用Spark Dataset.persist()复用缓存数据集?
跨EMR步骤复用Spark数据:persist()行不通,得这么做
先直接给你结论:这种情况下Spark会重新从S3读取所有数据,dataset.persist()完全帮不到跨EMR步骤共享数据的需求。原因和替代方案我给你拆解清楚:
为什么persist()跨EMR步骤无效?
- EMR步骤是完全独立的Spark应用:每个步骤对应一个单独的SparkSession/SparkContext,运行在独立的JVM进程中。当一个步骤执行完毕,整个应用的进程会被销毁,所有关联的内存资源(包括你persist的数据集)都会被系统回收,下个步骤根本访问不到这些缓存。
persist()的作用范围有限:它只能在当前Spark应用的生命周期内复用数据,缓存的数据存储在当前应用的Executor节点内存/磁盘上,跨应用(跨EMR步骤)完全隔离。
跨EMR步骤复用数据的正确做法
方案1:将中间数据落地到S3(推荐)
把需要复用的数据写入S3的中间路径,后续步骤直接从这个路径读取。建议用列式存储格式(比如Parquet、ORC),比text格式的读写效率高很多:
- 第一个EMR步骤:
val dataset = sqlContext.read.textFile("s3a://path/to/folder") // 如果有预处理逻辑可以在这里执行 dataset.write.parquet("s3a://path/to/intermediate-data") - 第二个EMR步骤:
val newDataset = sqlContext.read.parquet("s3a://path/to/intermediate-data")
方案2:合并任务到同一个EMR步骤
如果两个任务的逻辑可以串联,就把它们放到同一个Spark应用里,作为不同的阶段。这样persist()就能生效,不用重复读S3:
val dataset = sqlContext.read.textFile("s3a://path/to/folder") dataset.persist() // 缓存数据,供后续阶段复用 // 第一个任务逻辑 dataset.xxx() // 第二个任务逻辑,直接复用缓存的dataset dataset.xxx()
内容的提问来源于stack exchange,提问作者Abhay Dubey
相关产品推荐
相关产品推荐

