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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 02:13:18