【问题标题】:Read / Write Parquet files without reading into memory (using Python)读/写 Parquet 文件而不读入内存(使用 Python)
【发布时间】:2021-10-19 12:08:57
【问题描述】:

我查看了希望满足我需求的标准文档(Apache ArrowPandas),但我似乎无法弄清楚。

我最了解Python,所以我想用Python,但这不是一个严格的要求。

问题

我需要将 Parquet 文件从一个位置(一个 URL)移动到另一个位置(一个 Azure 存储帐户,在本例中使用 Azure 机器学习平台,但这与我的问题无关)。

这些文件太大而不能简单地执行pd.read_parquet("https://my-file-location.parquet"),因为这会将整个内容读入一个对象。

期待

我认为必须有一种简单 的方法来创建文件对象并逐行流式传输该对象——或者可能是逐列块。类似的东西

import pyarrow.parquet as pq

with pq.open("https://my-file-location.parquet") as read_file_handle:
    with pq.open("https://my-azure-storage-account/my-file.parquet", "write") as write_filehandle:
        for next_line in read_file_handle{
            write_file_handle.append(next_line)

我知道这会有点不同,因为 Parquet 主要是为了以柱状方式访问。也许我会传递某种配置对象来指定感兴趣的列,或者可以在一个块中抓取多少行或类似的东西。

但关键的期望是 有一种方法可以访问 parquet 文件而不会将其全部加载到内存中。我该怎么做?

FWIW,我确实尝试使用 Python 的标准 open 函数,但我不确定如何将 open 与 URL 位置和字节流一起使用。如果可以仅通过 open 执行此操作并跳过任何特定于 Parquet 的内容,那也可以。

更新

一些 cmets 建议使用类似 bash 的脚本,例如 here。如果没有别的,我可以使用它,但这并不理想,因为:

  • 我宁愿将这一切保存在一个完整的语言 SDK 中,无论是 Python、Go 还是其他任何东西。如果解决方案移动到带有管道的 bash 脚本中,则需要外部调用,因为最终解决方案不会完全由 bash、Powershell 或任何脚本语言编写。
  • 我真的很想利用 Parquet 本身的一些好处。正如我在下面的评论中提到的,Parquet 是列式存储。因此,如果我有一个 11 亿行和 100 列的“数据框”,但我只关心 3 列,我希望能够只下载这 3 列,从而节省大量时间和金钱。

【问题讨论】:

  • 这能回答你的问题吗? downloading a file from Internet into S3 bucket
  • 好吧,我没有使用 S3。而且,更一般地说,我不想依赖另一个可能广为人知但仍未像 HTTPS 那样完全通用的文件系统。所以我不想使用 S3、Databricks 的 DBFS 或 Azure 的 DFS 等。
  • 另外,需要明确的是,这个问题专门针对 Parquet。我假设有两件事:(1)如果我把它当作一个纯二进制文件并以某种方式流式传输它,这应该可以正常工作。所以@BeChillerToo 提到的curl 答案应该有效。 (2) 因为它是 Parquet,所以我应该具有能够在流式传输时即时进行一些简单处理的优势。例如,能够抓取感兴趣的列的子集而不是整个对象来存储。我希望这实际上可以根据要求完成,因此避免了不必要的网络 I/O。
  • 如果你想在运行中进行处理,这会大大改变微积分;将其视为纯二进制可能会更快,甚至可能更快,但仅适用于没有处理的直接副本

标签: python io parquet


【解决方案1】:

这是可能的,但需要做一些工作,因为除了柱状 Parquet 之外,还需要架构。

大致的工作流程是:

  1. 打开parquet file 阅读。

  2. 然后使用iter_batches 以增量方式回读行块(您也可以传递要从文件中读取的特定列以节省 IO/CPU)。

  3. 然后您可以将每个pa.RecordBatchiter_batches 进一步转换。完成第一批转换后,您可以获取其schema 并创建一个新的ParquetWriter

  4. 对于每个转换后的批处理调用 write_table。您必须先将其转换为pa.Table

  5. 关闭文件。

Parquet 需要随机访问,因此不能轻松地从 URI 流式传输(如果您通过 HTTP FSSpec 打开文件,pyarrow 应该支持它)但我认为您可能会在写入时被阻止。

【讨论】:

  • 另请参阅stackoverflow.com/questions/63891231/… 批量调整大小对于管理内存可能很重要。
  • 谢谢你,@Micah!是的,我知道这会很棘手,并且需要一个模式。我发现在 SO 上,如果我不真的知道我想问什么,有时写较短的问题比写较长的问题更好。 ;) 不管怎样,iter_batches 正是我正在寻找的。没看到就觉得傻。我会努力把它落实到位。
【解决方案2】:

请注意,我没有具体说明如何在远程服务器端使用批处理。
我的解决方案是:使用 pyarrow.NativeFile 将批处理写入缓冲区,然后使用 pyarrow.ipc.RecordBatchFileReader 读取缓冲区

我创建了这 2 个类来帮助您进行流式处理

import asyncio
from pyarrow.parquet import ParquetFile


class ParquetFileStreamer:
    """
    Attributes:
        ip_address: ip address of the distant server
        port: listening port of the distant server
        n_bytes: -1 means read whole batch
        file_source: pathlib.Path, pyarrow.NativeFile, or file-like object
        batch_size: default = 65536
        columns: list of the columns you wish to select (if None selects all)

    Example:
        >>> pfs = ParquetFileStreamer
        >>> class MyStreamer(ParquetFileStreamer)
                file_source = '/usr/fromage/camembert.parquet
                columns = ['name', 'price']
        >>> MyStreamer.start_stream()
    """
    ip_address = '192.168.1.1'
    port = 80
    n_bytes = -1

    file_source: str
    batch_size = 65536
    columns = []

    @classmethod
    def start_stream(cls):
        for batch in cls._open_parquet():
            asyncio.run(cls._stream_parquet(batch))

    @classmethod
    def _open_parquet(cls):
        return ParquetFile(cls.file_source).iter_batches(
            batch_size=cls.batch_size,
            columns=cls.columns
        )

    @classmethod
    async def _stream_parquet(cls, batch):
        reader, writer = await asyncio.open_connection(cls.ip_address, cls.port)
        writer.write(batch)
        await writer.drain()
        await reader.read()
        writer.close()
        await writer.wait_closed()


class ParquetFileReceiver:
    """
    Attributes: \n
        port: specify the port \n
        n_bytes: -1 reads all the batch
    Example:
        >>> pfr = ParquetFileReceiver
        >>> asyncio.run(pfr.server())
    """
    port = 80
    n_bytes = -1

    @classmethod
    async def handle_stream(cls, reader, writer):
        data = await reader.read(cls.n_bytes)
        batch = data.decode()
        print(batch)

    @classmethod
    async def server(cls):
        server = await asyncio.start_server(cls.handle_stream, port=cls.port)
        async with server:
            await server.serve_forever()

【讨论】:

    猜你喜欢
    • 2018-01-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-04-02
    • 2021-05-18
    • 1970-01-01
    相关资源
    最近更新 更多