写入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
相关产品推荐
相关产品推荐

