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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 02:55:14