使用Dask转换Avro到CSV时遭遇ValueError: read of closed file错误求助
处理大Avro文件时出现ValueError: read of closed file的解决方法
我需要处理一个超过2GB的Avro文件,想用Dask把它转成CSV格式(转JSON、Parquet也都失败了)。之前用类似方法实现CSV转JSON是正常的,但执行下面的代码时抛出了ValueError: read of closed file错误,附上代码和错误栈:
原代码
import dask.bag as db class Converter(): def __init__(self,input,output): """Converter constructor""" self.input = input self.output = output @staticmethod def large_file_reader(file_path: str): """ File reader """ temp_data = db.read_avro(file_path) data = temp_data.to_dataframe() # just to check able to by read properly print(data.head(6)) return data @staticmethod def large_file_writer(data, file_path: str) -> bool: """ File writer """ data.compute().to_csv(file_path, index=False) def large_file_processor(self): "Read then write" input_file_path =self.input output_file_path =self.output data = Converter.large_file_reader(input_file_path) Converter.large_file_writer(data=data, file_path=output_file_path) if __name__ == "__main__": c = Converter("/Users/csv_to_avro_new.avro", "/Users/test_avro_new.csv") c.large_file_processor()
错误栈
Traceback (most recent call last): File "/Users/PycharmProjects/ms--py/new.py", line 41, in <module> c.large_file_processor() File "/Users/PycharmProjects/ms--py/new.py", line 36, in large_file_processor Converter.large_file_writer(data=data, file_path=output_file_path) File "/Users/PycharmProjects/ms--py/new.py", line 28, in large_file_writer data.compute().to_csv(file_path, index=False) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/base.py", line 315, in compute (result,) = compute(self, traverse=False, **kwargs) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/base.py", line 600, in compute results = schedule(dsk, keys, **kwargs) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/threaded.py", line 89, in get results = get_async( File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/local.py", line 511, in get_async raise_exception(exc, tb) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/local.py", line 319, in reraise raise exc File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/local.py", line 224, in execute_task result = _execute_task(task, data) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task return func(*(_execute_task(a, cache) for a in args)) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/core.py", line 119, in <genexpr> return func(*(_execute_task(a, cache) for a in args)) File "/Users/adavsandeep/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/core.py", line 119, in _execute_task return func(*(_execute_task(a, cache) for a in args)) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/dask/bag/avro.py", line 150, in read_chunk chunk = read_block(f, off, l, head["sync"]) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/fsspec/utils.py", line 244, in read_block found_start_delim = seek_delimiter(f, delimiter, 2**16) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/fsspec/utils.py", line 187, in seek_delimiter current = file.read(blocksize) File "/Users/PycharmProjects/data-ingest/lib/python3.10/site-packages/fsspec/implementations/local.py", line 337, in read return self.f.read(*args, **kwargs) ValueError: read of closed file Process finished with exit code 1
问题原因
问题出在large_file_reader里的print(data.head(6))调用:
- Dask是延迟计算框架,
head()会触发一次实际的文件读取操作,读取完成后会关闭文件句柄 - 但返回的Dask DataFrame对象还保留着对原文件的读取任务引用,后续调用
compute()时,尝试再次读取已关闭的文件,就会抛出"read of closed file"错误
修复方案
1. 移除提前触发读取的调试代码
删掉print(data.head(6)),避免提前关闭文件句柄。
2. 使用Dask原生写入方法,避免全量加载到内存
原代码中data.compute().to_csv()会把整个大文件加载到内存转成Pandas DataFrame,不仅容易内存溢出,也是导致文件句柄问题的间接原因。改用Dask DataFrame的to_csv方法,保持分块处理的特性。
修改后的代码
import dask.bag as db class Converter(): def __init__(self, input, output): """Converter constructor""" self.input = input self.output = output @staticmethod def large_file_reader(file_path: str): """ File reader """ temp_data = db.read_avro(file_path) data = temp_data.to_dataframe() # 移除head()调用,避免提前触发文件读取 return data @staticmethod def large_file_writer(data, file_path: str) -> bool: """ File writer """ # 使用Dask DataFrame的to_csv方法,无需先compute() # single_file=True确保输出单个CSV文件,默认会生成多个分块文件 data.to_csv(file_path, index=False, single_file=True) return True def large_file_processor(self): "Read then write" input_file_path = self.input output_file_path = self.output data = Converter.large_file_reader(input_file_path) Converter.large_file_writer(data=data, file_path=output_file_path) if __name__ == "__main__": c = Converter("/Users/csv_to_avro_new.avro", "/Users/test_avro_new.csv") c.large_file_processor()
额外说明
如果需要验证Avro文件内容,可以用以下方式读取一小部分数据,避免触发全文件读取:
# 在large_file_reader中添加验证逻辑(可选) sample_data = data.sample(frac=0.001).compute() # 读取0.1%的数据 print(sample_data.head())
内容的提问来源于stack exchange,提问作者sandeep_a
相关产品推荐
相关产品推荐

