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

如何修复因写入中断损坏的Parquet文件?

问题

我有一个Parquet文件,它是通过如下循环持续写入生成的:

def process_data(self):
    #... other code ...

    with pq.ParquetWriter(self.destination_file, schema) as writer:
        with tqdm(total=total_rows, desc="Processing nodes") as pbar:
            for i in range(0, total_rows, self.batch_size):
                # ... processing code ...

                # Create a table from the batched data
                batch_table = pa.Table.from_arrays(
                    [
                        pa.array(node_ids),
                        pa.array(mut_positions),
                        pa.array(new_6mers),
                        pa.array(context_embeddings),
                        pa.array(nonmutation_contexts),
                    ],
                    schema=schema
                )

                # Write the batch table
                writer.write_table(batch_table)

                # ...

                pbar.update(len(batch_indices))

由于电脑中途意外关机,写入循环被强制中断。现在尝试使用pq.read_table读取该文件时,出现如下错误:

pyarrow.lib.ArrowInvalid: Error creating dataset. Could not read schema from 'data/processed/data_with_embeddings.parquet'. Is this a 'parquet' file?: Could not open Parquet input source 'data/processed/data_with_embeddings.parquet': Parquet magic bytes not found in footer. Either the file is corrupted or this is not a parquet file.

希望找到能保留大部分数据、仅丢失少量行的修复方案。

可行修复方案

Parquet文件损坏的核心原因是意外中断导致文件末尾的Footer(包含元数据、Schema等关键信息)未写入,但已经写入的Batch数据大概率是完整的。以下是具体修复步骤:

步骤1:定位文件中最后一个完整的Batch

Parquet的每个数据块末尾都带有PAR1魔法字节,我们可以通过扫描文件找到最后一个有效PAR1的位置,以此确定截断点。运行以下Python脚本:

import os

def find_last_valid_parquet_batch(file_path):
    par1_magic = b'PAR1'
    chunk_size = 4096
    file_size = os.path.getsize(file_path)
    
    # 从文件末尾往前扫描,减少IO次数
    with open(file_path, 'rb') as f:
        for offset in range(file_size, 0, -chunk_size):
            start_pos = max(0, offset - chunk_size)
            f.seek(start_pos)
            chunk_data = f.read(chunk_size)
            # 查找当前块中最后出现的PAR1标记
            magic_pos = chunk_data.rfind(par1_magic)
            if magic_pos != -1:
                # 返回最后一个PAR1标记的结束位置
                return start_pos + magic_pos + len(par1_magic)
    return -1

target_file = 'data/processed/data_with_embeddings.parquet'
last_valid_end = find_last_valid_parquet_batch(target_file)

步骤2:截断文件到有效位置

如果找到了有效位置,就将文件截断到该位置,去除末尾不完整的无效数据:

if last_valid_end != -1:
    with open(target_file, 'rb+') as f:
        f.truncate(last_valid_end)
    print(f"文件已修复,截断至{last_valid_end}字节位置")
else:
    print("未找到有效数据块标记,文件可能完全损坏")

步骤3:读取并验证修复后的文件

尝试用pyarrow读取修复后的文件,若成功则另存为新文件避免再次损坏:

import pyarrow.parquet as pq
import pyarrow as pa

try:
    repaired_table = pq.read_table(target_file)
    print(f"成功读取{repaired_table.num_rows}行有效数据")
    # 保存修复后的文件
    pq.write_table(repaired_table, 'data/processed/repaired_data.parquet')
except Exception as e:
    print(f"读取失败:{str(e)}")

备选方案:逐Batch读取容错

如果截断后仍无法读取,可使用低级API逐Batch读取,遇到损坏块时停止:

import pyarrow.parquet as pq
from pyarrow import parquet as pq_low

batches = []
schema = None

try:
    with open(target_file, 'rb') as f:
        parquet_reader = pq_low.ParquetFile(f)
        for batch in parquet_reader.iter_batches():
            batches.append(batch)
            if schema is None:
                schema = batch.schema
except Exception as e:
    print(f"读取到损坏块,已停止:{str(e)}")

if batches:
    valid_table = pa.Table.from_batches(batches, schema=schema)
    print(f"成功提取{valid_table.num_rows}行数据")
    pq.write_table(valid_table, 'data/processed/repaired_data.parquet')
else:
    print("未提取到任何有效数据")

内容的提问来源于stack exchange,提问作者dev-mirzabicer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 17:55:16