Spark未缓存却复用enrichment DataFrame,首次运行临时视图为空求助
问题分析与解决方案
核心原因
Spark的DataFrame采用惰性求值机制:你初始化enrichment变量时,只是定义了从parquet路径读取数据的逻辑,并没有实际执行读取操作。当你通过API生成新数据并追加到parquet路径后,之前定义的enrichment DataFrame并不会自动同步最新数据——首次执行join时,实际读取的还是初始的空数据集,所以临时视图为空;第二次运行时,enrichment重新读取了已经包含新数据的parquet路径,因此能正常关联并生成视图。
解决方法
必须在追加新数据后,重新读取enrichment存储路径,用最新的数据集和source进行关联,不能复用之前定义的enrichment变量。
修改后的代码示例
调整for循环中最后一步的逻辑,重新读取最新的enrichment数据:
for { entries <- IO { /* 找出所有新增的init条目逻辑 */ } completions <- enrichmentClient.enrich(apiKey, entries) _ <- IO { spark .createDataFrame(completions) .write .mode(SaveMode.Append) .parquet(enrichmentsLocation) } _ <- IO { // 关键:重新读取最新的enrichment数据 val updatedEnrichment = spark.read.schema(schema).parquet(enrichmentsLocation) source.join(updatedEnrichment, source(initText) === updatedEnrichment(initText)) .createOrReplaceTempView(sourceTableViewName) } } yield ()
额外注意事项
- 不要提前定义
enrichment变量并全程复用,Spark DataFrame的元数据不会自动感知底层存储的内容变化; - 如果后续有频繁读写需求,可考虑对
updatedEnrichment进行缓存,但每次数据更新后需先调用unpersist()清除旧缓存,再重新缓存新数据。
内容的提问来源于stack exchange,提问作者Mykeo
相关产品推荐
相关产品推荐

