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

写入gzip缓冲区时末尾记录丢失的问题求助

gzip写入时末尾记录丢失的问题排查与解决

问题描述

写入gzip缓冲区时出现末尾记录丢失:原始数据共4598行,但写入后内存缓冲区、临时文件及S3中均只有4587条记录。已调用gz_file.flush(),但问题仍存在。

相关代码

print(f'Data has {len(data)} rows')

try:
    with gzip.GzipFile(mode="w", fileobj=gz_buffer) as gz_file:
        writer = csv.DictWriter(TextIOWrapper(gz_file, "utf8"), fieldnames=data[0].keys())
        writer.writeheader()
        for i, row in enumerate(data):
            writer.writerow(row)
        print(f"{len(data)} records written")  # Log each record
        gz_file.flush() 
except Exception as e:
    print(f"Error writing record {e}")
    
    raise e

gz_buffer.seek(0)  # Important: seek back to the start of the buffer
with gzip.GzipFile(mode="r", fileobj=gz_buffer) as gz_file:
    reader = csv.DictReader(TextIOWrapper(gz_file, "utf8"))
    record_count = sum(1 for row in reader)  # Count records

print(f"Number of records in the buffer: {record_count}")

import tempfile

try:
    # Create a temporary file
    with tempfile.NamedTemporaryFile(mode="w+b", suffix=".gz", delete=False) as temp_file:
        with gzip.GzipFile(mode="w", fileobj=temp_file) as gz_file:
            writer = csv.DictWriter(TextIOWrapper(gz_file, "utf8"), fieldnames=data[0].keys())
            writer.writeheader()
            for i, row in enumerate(data):
                writer.writerow(row)
                #print(f"Record {i} written")
            gz_file.close()

        temp_file_path = temp_file.name
        print(f"Temporary file created at: {temp_file_path}")

    # Read from the temporary file
    with gzip.open(temp_file_path, 'rt') as gz_file:
        reader = csv.DictReader(gz_file)
        record_count = sum(1 for row in reader)

    print(f"Number of records in the temporary file: {record_count}")

执行输出

Data has 4598 rows
4598 records written
Number of records in the buffer: 4587
Temporary file created at: /var/folders/v8/hs0gswcx4459zbsz8zskbkz40000gr/T/tmpxf9nlmt6.gz
Number of records in the temporary file: 4587

问题原因

你只刷新了gzip.GzipFile对象的缓冲区,但上层的TextIOWrapper还有自己的内部缓冲区,csv写入的数据暂时滞留在这个缓冲区中,没有真正传递到底层的gz_buffer或临时文件里。手动调用gz_file.flush()无法触发TextIOWrapper的刷新操作,导致末尾部分记录未被写入。

另外,临时文件代码中手动调用gz_file.close()是多余的——with语句会自动处理文件关闭,手动关闭可能导致缓冲区未完全刷新就终止写入。

解决方法

方法一:使用write_through=True(Python 3.7+推荐)

创建TextIOWrapper时设置write_through=True,让所有写入操作直接透传到底层文件对象,跳过TextIOWrapper的缓冲区:

print(f'Data has {len(data)} rows')

try:
    with gzip.GzipFile(mode="w", fileobj=gz_buffer) as gz_file:
        # 开启write_through,直接透传写入
        writer = csv.DictWriter(TextIOWrapper(gz_file, "utf8", write_through=True), fieldnames=data[0].keys())
        writer.writeheader()
        for i, row in enumerate(data):
            writer.writerow(row)
        print(f"{len(data)} records written")
except Exception as e:
    print(f"Error writing record {e}")
    raise e

gz_buffer.seek(0)
with gzip.GzipFile(mode="r", fileobj=gz_buffer) as gz_file:
    reader = csv.DictReader(TextIOWrapper(gz_file, "utf8"))
    record_count = sum(1 for row in reader)

print(f"Number of records in the buffer: {record_count}")

方法二:手动刷新TextIOWrapper

如果无法使用write_through,可以显式获取TextIOWrapper对象并调用flush(),确保缓冲区数据全部写入底层:

print(f'Data has {len(data)} rows')

try:
    with gzip.GzipFile(mode="w", fileobj=gz_buffer) as gz_file:
        text_wrapper = TextIOWrapper(gz_file, "utf8")
        writer = csv.DictWriter(text_wrapper, fieldnames=data[0].keys())
        writer.writeheader()
        for i, row in enumerate(data):
            writer.writerow(row)
        print(f"{len(data)} records written")
        text_wrapper.flush()  # 先刷新TextIOWrapper缓冲区
        gz_file.flush()       # 再刷新gzip缓冲区
except Exception as e:
    print(f"Error writing record {e}")
    raise e

gz_buffer.seek(0)
with gzip.GzipFile(mode="r", fileobj=gz_buffer) as gz_file:
    reader = csv.DictReader(TextIOWrapper(gz_file, "utf8"))
    record_count = sum(1 for row in reader)

print(f"Number of records in the buffer: {record_count}")

临时文件代码修正

去掉手动gz_file.close(),并添加write_through=True:

import tempfile

try:
    with tempfile.NamedTemporaryFile(mode="w+b", suffix=".gz", delete=False) as temp_file:
        with gzip.GzipFile(mode="w", fileobj=temp_file) as gz_file:
            writer = csv.DictWriter(TextIOWrapper(gz_file, "utf8", write_through=True), fieldnames=data[0].keys())
            writer.writeheader()
            for i, row in enumerate(data):
                writer.writerow(row)
        temp_file_path = temp_file.name
        print(f"Temporary file created at: {temp_file_path}")

    with gzip.open(temp_file_path, 'rt') as gz_file:
        reader = csv.DictReader(gz_file)
        record_count = sum(1 for row in reader)

    print(f"Number of records in the temporary file: {record_count}")

核心总结

确保csv写入的数据完全传递到gzip文件的关键是处理TextIOWrapper的缓冲区,要么通过write_through=True跳过缓冲区,要么显式刷新TextIOWrapper。with语句会自动处理文件关闭,无需手动调用close(),避免提前终止写入流程。

内容的提问来源于stack exchange,提问作者L Xandor

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 15:28:18