读取Parquet文件后更新失败:refreshByPath无效的解决方法
Spark Parquet 状态持久化读写问题解决方案
问题场景
需要将平均值、斜率等状态数据持久化到Parquet文件,支持批处理作业内或跨作业的访问与更新。测试代码如下:
class TestStateManagement: @pytest.fixture(autouse=True) def init_test(self, spark_unit_test_fixture, tmp_path): self.spark = spark_unit_test_fixture self.path = tmp_path def test_double_safe_parquet(self): schema = StructType([StructField("time", IntegerType()), StructField("value", DoubleType())]) dataframe = self.spark.createDataFrame([(1, 1.0), (2, 2.0), (3, 3.0)], schema) parquet = os.path.join(self.path, "test.parquet") dataframe.write.mode("overwrite").parquet(parquet) dataframe = self.spark.read.parquet(parquet) # 此处引发后续问题 assert dataframe.count() == 3 dataframe_add = self.spark.createDataFrame([(4, 4.0), (5, 5.0), (6, 6.0)], schema) self.spark.catalog.refreshByPath(parquet) dataframe.union(dataframe_add).write.mode("overwrite").parquet(parquet)
错误信息
执行测试时触发文件不存在异常:
E Caused by: org.apache.spark.SparkFileNotFoundException: File file:/....snappy.parquet does not exist E It is possible the underlying files have been updated. You can explicitly invalidate the cache in Spark by running 'REFRESH TABLE tableName' command in SQL or by recreating the Dataset/DataFrame involved.
问题原因
Spark读取Parquet路径后会缓存文件的元数据信息。当通过overwrite模式写入新数据时,原Parquet文件会被删除,但之前创建的dataframe对象仍持有旧文件的引用。合并写入时,Spark尝试访问已被删除的旧文件,从而抛出异常。refreshByPath无法解决该问题,因为它仅作用于Catalog中注册的表,而直接读取路径生成的DataFrame元数据缓存不受Catalog管理。
解决方案
方案1:重新读取最新数据再合并
放弃复用之前的dataframe对象,在合并操作前重新读取Parquet文件,确保使用最新的文件元数据:
def test_double_safe_parquet(self): schema = StructType([StructField("time", IntegerType()), StructField("value", DoubleType())]) dataframe = self.spark.createDataFrame([(1, 1.0), (2, 2.0), (3, 3.0)], schema) parquet = os.path.join(self.path, "test.parquet") dataframe.write.mode("overwrite").parquet(parquet) # 验证数据时读取一次 assert self.spark.read.parquet(parquet).count() == 3 dataframe_add = self.spark.createDataFrame([(4, 4.0), (5, 5.0), (6, 6.0)], schema) # 合并前重新读取最新数据 latest_df = self.spark.read.parquet(parquet) latest_df.union(dataframe_add).write.mode("overwrite").parquet(parquet)
方案2:使用临时表配合REFRESH TABLE
将Parquet路径注册为临时表,通过Catalog管理元数据,此时REFRESH TABLE可以有效清除缓存:
def test_double_safe_parquet(self): schema = StructType([StructField("time", IntegerType()), StructField("value", DoubleType())]) dataframe = self.spark.createDataFrame([(1, 1.0), (2, 2.0), (3, 3.0)], schema) parquet = os.path.join(self.path, "test.parquet") dataframe.write.mode("overwrite").parquet(parquet) # 注册为临时表 self.spark.sql(f"CREATE OR REPLACE TEMP VIEW state_view AS SELECT * FROM parquet.`{parquet}`") # 通过临时表验证数据 assert self.spark.table("state_view").count() == 3 dataframe_add = self.spark.createDataFrame([(4, 4.0), (5, 5.0), (6, 6.0)], schema) # 刷新临时表元数据 self.spark.sql("REFRESH TABLE state_view") # 从临时表读取最新数据合并 self.spark.table("state_view").union(dataframe_add).write.mode("overwrite").parquet(parquet)
方案3:避免持有旧DataFrame引用
如果必须保留中间读取的验证步骤,确保验证后不再使用该DataFrame对象,而是重新生成新的DataFrame进行后续操作。
内容的提问来源于stack exchange,提问作者dermoritz
相关产品推荐
相关产品推荐

