代码之家  ›  专栏  ›  技术社区  ›  suvayu

如何从多个进程(可能来自不同的语言)使用Apache Arrow IPC?

  •  0
  • suvayu  · 技术社区  · 3 年前

    我不知道从哪里开始,所以寻求一些指导。我正在寻找一种方法,在一个进程中创建一些数组/表,并从另一个进程访问(只读)。

    所以我创建了一个 pyarrow.Table 这样地:

    a1 = pa.array(list(range(3)))
    a2 = pa.array(["foo", "bar", "baz"])
    
    a1
    # <pyarrow.lib.Int64Array object at 0x7fd7c4510100>
    # [
    #   0,
    #   1,
    #   2
    # ]
    
    a2
    # <pyarrow.lib.StringArray object at 0x7fd7c5d6fa00>
    # [
    #   "foo",
    #   "bar",
    #   "baz"
    # ]
    
    tbl = pa.Table.from_arrays([a1, a2], names=["num", "name"])
    
    tbl
    # pyarrow.Table
    # num: int64
    # name: string
    # ----
    # num: [[0,1,2]]
    # name: [["foo","bar","baz"]]
    

    现在,我如何从不同的过程中解读这一点?我想我会用 multiprocessing.shared_memory.SharedMemory ,但这并不完全奏效:

    shm = shared_memory.SharedMemory(name='pa_test', create=True, size=tbl.nbytes)
    with pa.ipc.new_stream(shm.buf, tbl.schema) as out:
        for batch in tbl.to_batches():
            out.write(batch)
    
    # TypeError: Unable to read from object of type: <class 'memoryview'>
    

    我需要包一下吗 shm.buf 用什么?

    即使我能做到这一点,它似乎也很麻烦。我该如何以稳健的方式做到这一点?我需要像zmq这样的东西吗?

    但我不清楚这怎么是零拷贝。当我写记录批次时,这不是串行化吗?我错过了什么?

    在我的真实用例中,我也想和Julia谈谈,但也许这应该是一个单独的问题。

    附言:我已经经历了 docs ,它没有为我澄清这一部分。

    0 回复  |  直到 3 年前
        1
  •  9
  •   Will Jones    3 年前

    我需要用什么东西包shm.buf吗?

    是的,你可以使用 pa.py_buffer() 要包装它:

    size = calculate_ipc_size(table)
    shm = shared_memory.SharedMemory(create=True, name=name, size=size)
    
    stream = pa.FixedSizeBufferWriter(pa.py_buffer(shm.buf))
    with pa.RecordBatchStreamWriter(stream, table.schema) as writer:
       writer.write_table(table)
    

    此外,对于 size 您需要计算IPC输出的大小,它可能比 Table.nbytes 。您可以使用的功能是:

    def calculate_ipc_size(table: pa.Table) -> int:
        sink = pa.MockOutputStream()
        with pa.ipc.new_stream(sink, table.schema) as writer:
            writer.write_table(table)
        return sink.size()
    

    我该如何以稳健的方式做到这一点?

    还不确定这部分。根据我的经验,当其他进程重用缓冲区时,原始进程需要保持活力,但可能有办法绕过这一点。这可能与CPython中的此错误有关: https://bugs.python.org/issue38119

    但我不清楚这怎么是零拷贝。当我写记录批次时,这不是串行化吗?我错过了什么?

    你说得对 写 进入IPC缓冲器的箭头数据确实涉及拷贝。零拷贝部分是当其他进程 阅读 共享内存中的数据。Arrow表的列将引用IPC缓冲区的相关段,而不是副本。