PyArrow类型DataFrame存入共享内存后关闭报错求助
解决PyArrow类型DataFrame存入共享内存后的BufferError问题
问题现象
当使用PyArrow类型(如string[pyarrow])的DataFrame存入多进程共享内存时,关闭共享内存会出现两类异常:
- 调用
sm_put.close()后,读取端执行sm_get.close()直接触发BufferError: memoryview has 1 exported buffer,导致代码终止 - 仅关闭写入端共享内存,虽无即时错误,但一段时间后会弹出相同错误警告
经memory_profiler观察,删除PyArrow对象(如table、record_batch)后内存未即时释放,推测是PyArrow对象仍持有共享内存缓冲区引用,导致共享内存无法正常释放。
核心原因
PyArrow的RecordBatch、Table等对象会隐式持有底层内存缓冲区的引用,即使显式执行del删除对象,Python垃圾回收机制不会立即触发,导致共享内存的memoryview仍存在导出的buffer引用,无法正常释放。
修复方案
1. 显式触发垃圾回收
在清理PyArrow对象后,强制调用gc.collect(),立即回收未引用的对象,释放共享内存buffer的引用。
2. 用with语句管理资源
使用with语句自动管理PyArrow的StreamWriter、StreamReader等资源,避免手动关闭遗漏导致的引用残留。
3. 调整资源释放顺序
确保在关闭共享内存前,所有依赖共享内存buffer的PyArrow对象已被彻底清理,避免交叉引用。
修改后的代码示例
写入共享内存函数
import gc import sys import pyarrow as pa import pandas as pd from multiprocessing import SharedMemory def put_df(data, logTime, logSize, usingArrow=False): record_batch = None table = None try: if usingArrow: table = pa.Table.from_pandas(data).combine_chunks() record_batch = table.to_batches(max_chunksize=sys.maxsize)[0] else: record_batch = pa.RecordBatch.from_pandas(data) # 计算所需共享内存大小 mock_sink = pa.MockOutputStream() with pa.RecordBatchStreamWriter(mock_sink, record_batch.schema) as stream_writer: stream_writer.write_batch(record_batch) data_size = mock_sink.size() sm_put = SharedMemory(create=True, size=data_size) buffer = pa.py_buffer(sm_put.buf) # 写入共享内存 with pa.FixedSizeBufferWriter(buffer) as stream: with pa.RecordBatchStreamWriter(stream, record_batch.schema) as stream_writer: stream_writer.write_batch(record_batch) # 清理PyArrow对象 del record_batch if usingArrow: del table del buffer # 强制触发垃圾回收 gc.collect() sm_put.close() return sm_put.name finally: # 兜底清理,确保异常场景下资源释放 if 'record_batch' in locals(): del record_batch if 'table' in locals(): del table gc.collect()
读取共享内存函数
def get_df(sm_get_name, logTime, logSize, usingArrow=False): sm_get = SharedMemory(name=sm_get_name, create=False) buffer = None reader = None record_batch = None try: buffer = pa.BufferReader(sm_get.buf) reader = pa.RecordBatchStreamReader(buffer) record_batch = reader.read_next_batch() # 转换为DataFrame if usingArrow: data = record_batch.to_pandas(types_mapper=pd.ArrowDtype) else: data = record_batch.to_pandas() return data finally: # 清理所有PyArrow对象并触发GC if 'record_batch' in locals(): del record_batch if 'reader' in locals(): del reader if 'buffer' in locals(): del buffer gc.collect() # 最后关闭并删除共享内存 sm_get.close() sm_get.unlink()
额外建议
- 建议升级PyArrow至12.0.0及以上版本,旧版本存在内存引用管理的已知bug
- 避免在多进程中直接持有共享内存的buffer引用,确保所有操作完成后立即释放相关对象
内容的提问来源于stack exchange,提问作者hst01
相关产品推荐
相关产品推荐

