Spark中lazy val的DataFrame读取器重复执行S3读取问题求助
问题分析
你遇到的核心问题不是lazy val失效,而是**Reader的执行逻辑特性**导致的:
studentDataReader作为lazy val确实只会初始化一次Reader实例,但这个Reader的run方法每次执行时,都会重新调用传入的config => getStudentData(config)函数,每次调用都会生成新的DataFrame实例。- 你在
getStudentData中调用的cache()是Spark的懒加载缓存,只有当DataFrame被action操作触发时才会实际缓存数据。但因为每次run都会生成新的DataFrame,两次调用studentDataReader会产生两个独立的DataFrame,各自触发S3读取。
解决方案
以下是三种可行的解决思路,按推荐优先级排序:
方案一:缓存Reader的执行结果(最简洁)
在CommonProcessor中用线程安全的缓存结构,针对同一个Setting只执行一次读取逻辑,并强制触发Spark缓存:
trait CommonProcessor { val spark: SparkSession = SparkSession.getActiveSession.get import spark.implicits._ // 用ConcurrentHashMap缓存不同Setting对应的学生数据DF private val studentDataCache = new java.util.concurrent.ConcurrentHashMap[Setting, DataFrame]() lazy val studentDataReader: Reader[Setting, DataFrame] = Reader({ config => studentDataCache.computeIfAbsent(config, setting => { val studentData = getStudentData(setting) // 强制触发缓存:Spark cache是懒加载,必须执行action才会实际缓存数据 studentData.count() studentData }) }) def getStudentData(setting: Setting): DataFrame = { val studentData = getInput(setting, deltaPathToS3) //底层调用spark.read.format("delta") studentData.cache() } }
关键点:
computeIfAbsent保证同一个Setting只会调用一次getStudentDatacount()触发缓存,避免后续操作重复读取S3
方案二:复用DataFrame实例(避免多次调用Reader)
在StudentDataProcessor中先获取一次studentData,再传递给所有需要的逻辑,彻底避免两次调用studentDataReader:
trait StudentDataProcessor { self: CommonProcessor => import spark.implicits._ lazy val getResult: Reader[Setting, DataFrame] = { for { df <- getFinalData } yield { dropDuplicates(df) } } lazy val getFinalData: Reader[Setting, DataFrame] = { Reader({ config => // 只调用一次studentDataReader val studentData = studentDataReader.run(config) // 直接复用studentData生成关联数据 val studentTeacherData = mergeData(studentData, teacherDataReader.run(config)) doSomeProcess(studentData, studentTeacherData) }) } }
方案三:给Reader添加记忆化能力(如果可修改Reader trait)
如果能自定义Reader的实现,可以给它添加memoize方法,让run结果自动缓存:
trait Reader[Config, A] { def run(config: Config): A // 添加记忆化方法 def memoize: Reader[Config, A] = { val cache = new java.util.concurrent.ConcurrentHashMap[Config, A]() Reader({ config => cache.computeIfAbsent(config, run) }) } }
然后修改CommonProcessor中的studentDataReader:
lazy val studentDataReader: Reader[Setting, DataFrame] = Reader({ config => val studentData = getStudentData(config) studentData.count() // 触发缓存 studentData }).memoize
内容的提问来源于stack exchange,提问作者hsnsd
相关产品推荐
相关产品推荐

