如何修复因写入中断损坏的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
相关产品推荐
相关产品推荐

