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

Polars使用scan_ipc从共享内存Arrow IPC缓冲区懒加载数据失败求助

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.13 18:29:34