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

使用谓词从pyarrow中筛选行。拼花地板ParquetDataset

  •  0
  • kluu  · 技术社区  · 7 年前

    我有一个拼花地板数据集存储在s3上,我想从数据集中查询特定的行。我可以用 petastorm 但现在我只想用 pyarrow 。

    以下是我的尝试:

    import pyarrow.parquet as pq
    import s3fs
    
    fs = s3fs.S3FileSystem()
    
    dataset = pq.ParquetDataset(
        'analytics.xxx', 
        filesystem=fs, 
        validate_schema=False, 
        filters=[('event_name', '=', 'SomeEvent')]
    )
    
    df = dataset.read_pandas().to_pandas()
    

    但这会返回一个pandas数据帧,就好像过滤器不起作用一样,也就是说,我有多个值为 event_name .我是否有遗漏或误解的地方?我可以在获得熊猫数据帧后进行过滤,但我会使用比需要更多的内存空间。

    0 回复  |  直到 7 年前
        1
  •  19
  •   rjurney Sean Vieira    5 年前

    注意:我已将其扩展为Python和Parquet的全面指南 this post

    拼花地板格式分区

    为了使用过滤器,您需要使用分区以拼花格式存储数据。从许多拼花地板列和分区中加载几个拼花地板列和分区可以大大提高拼花地板与CSV的I/O性能。Parquet可以基于一个或多个字段的值对文件进行分区,并为嵌套值的唯一组合创建目录树,或者为一个分区列创建一组目录。这个 PySpark Parquet documentation 解释拼花地板如何工作得相当好。

    性别和国家的划分 this :

    path
    └── to
        └── table
            ├── gender=male
            │   ├── ...
            │   │
            │   ├── country=US
            │   │   └── data.parquet
            │   ├── country=CN
            │   │   └── data.parquet
            │   └── ...
    

    如果需要对数据进行进一步分区,还可以进行行组分区,但大多数工具只支持指定行组大小,您必须执行以下操作 key-->row group 查找你自己,这很难看(很高兴在另一个问题中回答这个问题)。

    使用Pandas编写分区

    您需要使用Parquet对数据进行分区,然后可以使用过滤器加载数据。您可以使用PyArrow、pandas或 Dask 或 PySpark 对于大型数据集。

    例如,要在pandas中写入分区,请执行以下操作:

    df.to_parquet(
        path='analytics.xxx', 
        engine='pyarrow',
        compression='snappy',
        columns=['col1', 'col5'],
        partition_cols=['event_name', 'event_category']
    )
    

    这会将文件按如下方式排列:

    analytics.xxx/event_name=SomeEvent/event_category=SomeCategory/part-0001.c000.snappy.parquet
    analytics.xxx/event_name=SomeEvent/event_category=OtherCategory/part-0001.c000.snappy.parquet
    analytics.xxx/event_name=OtherEvent/event_category=SomeCategory/part-0001.c000.snappy.parquet
    analytics.xxx/event_name=OtherEvent/event_category=OtherCategory/part-0001.c000.snappy.parquet
    

    在PyArrow中加载拼花地板分区

    要使用分区列按一个属性获取事件,请在列表中放置元组过滤器:

    import pyarrow.parquet as pq
    import s3fs
    
    fs = s3fs.S3FileSystem()
    
    dataset = pq.ParquetDataset(
        's3://analytics.xxx', 
        filesystem=fs, 
        validate_schema=False, 
        filters=[('event_name', '=', 'SomeEvent')]
    )
    df = dataset.to_table(
        columns=['col1', 'col5']
    ).to_pandas()
    

    使用逻辑AND进行筛选

    要获取具有两个或多个属性的事件,只需创建过滤器元组列表:

    import pyarrow.parquet as pq
    import s3fs
    
    fs = s3fs.S3FileSystem()
    
    dataset = pq.ParquetDataset(
        's3://analytics.xxx', 
        filesystem=fs, 
        validate_schema=False, 
        filters=[
            ('event_name',     '=', 'SomeEvent'),
            ('event_category', '=', 'SomeCategory')
        ]
    )
    df = dataset.to_table(
        columns=['col1', 'col5']
    ).to_pandas()
    

    使用逻辑OR筛选

    要使用或获取两个事件,需要将过滤器元组嵌套在它们自己的列表中:

    import pyarrow.parquet as pq
    import s3fs
    
    fs = s3fs.S3FileSystem()
    
    dataset = pq.ParquetDataset(
        's3://analytics.xxx', 
        filesystem=fs, 
        validate_schema=False, 
        filters=[
            [('event_name', '=', 'SomeEvent')],
            [('event_name', '=', 'OtherEvent')]
        ]
    )
    df = dataset.to_table(
        columns=['col1', 'col5']
    ).to_pandas()
    

    使用AWS Data Wrangler加载拼花地板分区

    正如前面提到的另一个答案,将数据过滤加载到数据所在的特定分区(本地或云中)的特定列的最简单方法是使用 awswrangler 单元如果您使用的是S3,请查看以下文档: awswrangler.s3.read_parquet() 和 awswrangler.s3.to_parquet() .过滤的工作原理与上述示例相同。

    import awswrangler as wr
    
    df = wr.s3.read_parquet(
        path="analytics.xxx",
        columns=["event_name"], 
        filters=[('event_name', '=', 'SomeEvent')]
    )
    

    加载拼花地板隔墙 pyarrow.parquet.read_table()

    如果您使用的是PyArrow,那么还可以使用 pyarrow。拼花地板read\u table() :

    import pyarrow.parquet as pq
    
    fp = pq.read_table(
        source='analytics.xxx',
        use_threads=True,
        columns=['some_event', 'some_category'],
        filters=[('event_name', '=', 'SomeEvent')]
    )
    df = fp.to_pandas()
    

    使用PySpark加载拼花地板分区

    最后,在PySpark中,您可以使用 pyspark.sql.DataFrameReader.read_parquet()

    import pyspark.sql.functions as F
    from pyspark.sql import SparkSession
    
    spark = SparkSession.builder.master("local[1]") \
                        .appName('Stack Overflow Example Parquet Column Load') \
                        .getOrCreate()
    
    # I automagically employ Parquet structure to load the selected columns and partitions
    df = spark.read.parquet('s3://analytics.xxx') \
              .select('event_name', 'event_category') \
              .filter(F.col('event_name') == 'SomeEvent')
    

    希望这能帮助您使用拼花地板:)

        2
  •  13
  •   Niklas B    6 年前

    对于任何从Google来到这里的人来说,当你阅读拼花地板文件时,你现在可以在PyArrow中过滤行。不管你是通过pandas还是pyarrow阅读。拼花地板

    从 documentation :

    过滤器 (列表[元组]或列表[列表[元组]]或无(默认)) 将从扫描的数据中删除与筛选器谓词不匹配的行。嵌入在嵌套目录结构中的分区键将被利用,以避免加载不包含匹配行的文件。如果use\u legacy\u dataset为True,则筛选器只能引用分区键,并且仅支持配置单元样式的目录结构。将use\u legacy\u dataset设置为False时,还支持文件级过滤和不同的分区方案。

    谓词以析取范式(DNF)表示,如[[('x','=',0),…],…]。DNF允许单列谓词的任意布尔逻辑组合。最内层的元组分别描述一个列谓词。内部谓词列表被解释为连词(AND),形成一个更具选择性的多列谓词。最后,最外层的列表将这些过滤器组合为析取(OR)。

    谓词也可以作为列表[元组]传递。这种形式被解释为单一连词。要在谓词中表示OR,必须使用(首选)List[List[Tuple]]表示法。

        3
  •  4
  •   joris    7 年前

    目前 filters 功能仅在文件级别实现,尚未在行级别实现。

    因此,如果数据集是嵌套层次结构中多个分区拼花文件的集合(此处描述的分区数据集类型: https://arrow.apache.org/docs/python/parquet.html#partitioned-datasets-multiple-files ),您可以使用 过滤器 参数仅读取文件的子集。
    但是,您还不能将其仅用于读取单个文件行组的子集(请参见 https://issues.apache.org/jira/browse/ARROW-1796 )。

    但是,如果您得到一条指定这样一个无效过滤器的错误消息,那就太好了。我为此打开了一个问题: https://issues.apache.org/jira/browse/ARROW-5572

        4
  •  4
  •   Vincent Claes    6 年前

    对于python 3.6+而言,AWS有一个名为AWS data wrangler的库,该库有助于Pandas/S3/Parquet之间的集成,并允许您过滤已分区的S3键。

    安装do;

    pip install awswrangler
    

    为了减少读取的数据,可以根据存储在s3上的拼花地板文件中的分区列过滤行。 从分区列中筛选行 event_name 具有值 "SomeEvent" 做

    对于awswrangler<1.0.0

    import awswrangler as wr
    
    df = wr.pandas.read_parquet(
             path="s3://my-bucket/my/path/to/parquet-file.parquet",
             columns=["event_name"], 
             filters=[('event_name', '=', 'SomeEvent')]
    )
    

    对于awswrangler>1.0.0 do;

    import awswrangler as wr
    
    df = wr.s3.read_parquet(
             path="s3://my-bucket/my/path/to/parquet-file.parquet",
             columns=["event_name"], 
             filters=[('event_name', '=', 'SomeEvent')]
    )