Polars使用scan_ipc从共享内存Arrow IPC缓冲区懒加载数据失败求助
问题描述
我最近在做一个从共享内存里通过Arrow IPC读写Polars数据的系统,eager模式下用read_ipc读取完全没问题,但换成懒加载的scan_ipc就一直报错,折腾了好半天都没解决。下面是我写的最小复现脚本:
# /// script # requires-python = "~=3.12.0" # dependencies = [ # "polars[pyarrow]==1.38.1", # ] # /// """Minimal working example for scan IPC bug.""" import logging from multiprocessing.shared_memory import SharedMemory import polars as pl import pyarrow as pa LOG = logging.getLogger(__name__) pl.show_versions() if __name__ == "__main__": # Create a toy dataset data = pl.DataFrame({"a": [1, 2, 3]}) with pa.MockOutputStream() as sink: data.write_ipc(sink) shm = SharedMemory(name="test-ipc-polars", create=True, size=sink.size()) with pa.FixedSizeBufferWriter(pa.py_buffer(shm.buf)) as stream: data.write_ipc(stream) del stream shm.close() # Read the dataset back existing_ = SharedMemory(shm.name, create=False) buf_ = pa.py_buffer(existing_.buf) # out = pl.read_ipc(buf_, use_pyarrow=True) # Works # print(out) try: res = pl.scan_ipc(buf_) except Exception as exc: LOG.error(exc, stack_info=True, exc_info=True) finally: del buf_ existing_.close() shm.unlink()
运行脚本后,会抛出以下错误:
_init_credential_provider_builder(): credential_provider_init = None async thread count: 4 polars-stream: updating graph state polars-stream: running in-memory-source in subgraph polars-stream: running io-sink[single-file[ipc]] in subgraph async upload_chunk_size: 67108864 async upload_concurrency: 8 io-sink[single-file[ipc]]: start_single_file_sink_pipeline: file_writer_starter: ipc, takeable_rows_provider: TakeableRowsProvider { max_size: NonZeroRowCountAndSize { num_rows: 122880, num_bytes: 18446744073709551615 }, byte_size_min_rows: 16384, allow_non_max_size: false }, upload_chunk_size: 67108864 polars-stream: done running graph phase polars-stream: updating graph state io-sink[single-file[ipc]]: Join on task_handle (recv PortState::Done) io-sink[single-file[ipc]]: Statistics: total_size: RowCountAndSize { num_rows: 3, num_bytes: 24 } _init_credential_provider_builder(): credential_provider_init = None polars-stream: updating graph state polars-stream: running io-sink[single-file[ipc]] in subgraph polars-stream: running in-memory-source in subgraph io-sink[single-file[ipc]]: start_single_file_sink_pipeline: file_writer_starter: ipc, takeable_rows_provider: TakeableRowsProvider { max_size: NonZeroRowCountAndSize { num_rows: 122880, num_bytes: 18446744073709551615 }, byte_size_min_rows: 16384, allow_non_max_size: false }, upload_chunk_size: 67108864 polars-stream: done running graph phase polars-stream: updating graph state io-sink[single-file[ipc]]: Join on task_handle (recv PortState::Done) io-sink[single-file[ipc]]: Statistics: total_size: RowCountAndSize { num_rows: 3, num_bytes: 24 } _init_credential_provider_builder(): credential_provider_init = None argument 'sources': Object does not have a .read() method. Traceback (most recent call last): File "/workspaces/github/mwe.py", line 36, in <module> res = pl.scan_ipc(buf_) ^^^^^^^^^^^^^^^^^ File "/app/.cache/uv/environments-v2/mwe-63942327ce1a05a8/lib/python3.12/site-packages/polars/_utils/deprecation.py", line 128, in wrapper return function(*args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^ File "/app/.cache/uv/environments-v2/mwe-63942327ce1a05a8/lib/python3.12/site-packages/polars/_utils/deprecation.py", line 128, in wrapper return function(*args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^ File "/app/.cache/uv/environments-v2/mwe-63942327ce1a05a8/lib/python3.12/site-packages/polars/io/ipc/functions.py", line 506, in scan_ipc pylf = PyLazyFrame.new_from_ipc( ^^^^^^^^^^^^^^^^^^^^^^^^^ TypeError: argument 'sources': Object does not have a .read() method. Stack (most recent call last): File "/workspaces/github/mwe.py", line 38, in <module> LOG.error(exc, stack_info=True, exc_info=True)
我特别希望能实现对这个共享内存缓冲区的懒加载,而不是一次性把数据全读到内存里。
可能的解决方案
1. 把PyArrow Buffer包装成类文件对象
pl.scan_ipc要求传入的数据源必须具备.read()方法(类似文件对象的接口),但直接传入pa.py_buffer是不满足这个要求的。我们可以用io.BytesIO把buffer包装成类文件对象:
import io # ... 前面的代码保持不变 ... existing_ = SharedMemory(shm.name, create=False) buf_ = pa.py_buffer(existing_.buf) # 将共享内存buffer转为类文件对象 file_like_buf = io.BytesIO(buf_.to_pybytes()) try: res = pl.scan_ipc(file_like_buf) # 验证懒加载效果 print(res.collect()) except Exception as exc: LOG.error(exc, stack_info=True, exc_info=True) finally: del buf_ existing_.close() shm.unlink()
不过要注意:这种方式会把共享内存里的数据复制到BytesIO的内存空间中,虽然能解决scan_ipc的调用问题,但会失去共享内存零拷贝的优势,如果你的数据量很大,这个方案可能不太理想。
2. 升级Polars版本
你当前使用的是Polars 1.38.1,可能这个版本的scan_ipc对内存缓冲区的支持还不完善。可以尝试升级到最新版本的Polars,看看官方是否已经修复了这个问题,或者新增了直接支持PyArrow Buffer作为scan_ipc输入的特性。
3. 通过PyArrow Dataset中转
另一种思路是先用PyArrow读取共享内存buffer并创建Dataset,再交给Polars扫描:
# ... 前面的代码保持不变 ... existing_ = SharedMemory(shm.name, create=False) buf_ = pa.py_buffer(existing_.buf) # 用PyArrow读取buffer并生成Dataset pa_table = pa.ipc.open_stream(buf_).read_all() pa_dataset = pa.dataset.dataset(pa_table) try: res = pl.scan_pyarrow_dataset(pa_dataset) print(res.collect()) except Exception as exc: LOG.error(exc, stack_info=True, exc_info=True) finally: del buf_ existing_.close() shm.unlink()
不过这个方案本质上还是先把数据加载到PyArrow Table中(属于eager加载),只是后续可以用Polars的懒加载接口处理数据,如果你的目标是完全的懒加载(即仅在需要时才读取共享内存中的数据),这个方法可能不符合预期。
总结
目前pl.scan_ipc对无.read()方法的内存缓冲区支持有限,最直接的临时解决方案是把buffer包装成类文件对象,但会牺牲零拷贝特性。长期来看,建议关注Polars的版本更新,或者在Polars的GitHub仓库提交issue,请求官方支持直接从PyArrow Buffer懒加载数据的功能。
备注:内容来源于stack exchange,提问作者akgcodes

