You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

读取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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.20 07:32:19