【问题标题】:Read multiple parquet files with selected columns into one Pandas dataframe将具有选定列的多个镶木地板文件读入一个 Pandas 数据帧
【发布时间】:2022-01-17 02:09:02
【问题描述】:

我正在尝试将具有选定列的多个镶木地板文件读入一个 Pandas 数据帧。这意味着镶木地板文件不共享所有列。我试图在pd.read_parquet() 中添加一个filter() 参数,但它似乎在多文件读取中不起作用。我怎样才能做到这一点?

from pathlib import Path
import pandas as pd

data_dir = Path('dir/to/parquet/files')
full_df = pd.concat(
    pd.read_parquet(parquet_file)
    for parquet_file in data_dir.glob('*.parquet')
)


full_df = pd.concat(
    pd.read_parquet(parquet_file, filters=[('name', 'address', 'email')])
    for parquet_file in data_dir.glob('*.parquet')
)

【问题讨论】:

  • 我认为您的意思是使用columns 而不是filterscolumns 用于提取特定列,filters 用于过滤带有条件的数据(即:[('name', '>', 'Mary')])。另外,对我来说,将目录路径传递给read_parquet 读取所有文件而不使用 for 循环。我正在使用熊猫==1.2.4。

标签: python pandas pyarrow


【解决方案1】:

很好地支持从多个文件中读取。但是,如果您的架构不同,那就有点棘手了。 Pyarrow 当前默认使用它在数据集中找到的第一个文件的模式。这是为了避免检查大型数据集中每个文件的架构的前期成本。

Arrow-C++ 有 the capability 覆盖它并扫描每个文件,但这还没有在 pyarrow 中公开。但是,如果您提前知道统一模式,则可以提供它,您将获得所需的行为。您将需要直接使用数据集模块来执行此操作,因为指定架构不是 pyarrow.parquet.read_table 的一部分(这是由 pandas.read_parquet 调用的)。

import pyarrow as pa
import pyarrow.dataset as ds
import pyarrow.parquet as pq
import pandas as pd

import tempfile

tab1 = pa.Table.from_pydict({'a': [1, 2, 3], 'b': ['a', 'b', 'c']})
tab2 = pa.Table.from_pydict({'b': ['a', 'b', 'c'], 'c': [True, False, True]})

unified_schema = pa.unify_schemas([tab1.schema, tab2.schema])

with tempfile.TemporaryDirectory() as dataset_dir:

    pq.write_table(tab1, f'{dataset_dir}/one.parquet')
    pq.write_table(tab2, f'{dataset_dir}/two.parquet')

    print('Basic read of directory will use schema from first file')
    print(pd.read_parquet(dataset_dir))

    print()
    print('You can specify the unified schema if you know it')
    dataset = ds.dataset(dataset_dir, schema=unified_schema)
    print(dataset.to_table().to_pandas())

    print()
    print('The columns option will limit which columns are returned from read_parquet')
    print(pd.read_parquet(dataset_dir, columns=['b']))

    print()
    print('The columns option can be used when specifying a schema as well')
    dataset = ds.dataset(dataset_dir, schema=unified_schema)
    print(dataset.to_table(columns=['b', 'c']).to_pandas())

如果您不提前知道统一架构,您可以自己检查所有文件来创建它:

# You could also use glob here or whatever tool you want to
# get the list of files in your dataset
dataset = ds.dataset(dataset_dir)
schemas = [pq.read_schema(dataset_file) for dataset_file in dataset.files]
print(pa.unify_schemas(schemas))

由于这可能很昂贵(尤其是在使用远程文件系统时),您可能希望将统一架构保存在自己的文件中(保存一个 parquet 文件或 0 个批次的 Arrow IPC 文件通常就足够了)而不是重新计算它每次。

【讨论】:

    猜你喜欢
    • 2020-03-28
    • 1970-01-01
    • 2018-12-22
    • 1970-01-01
    • 2018-06-24
    • 2019-02-11
    • 2020-03-23
    • 2020-04-02
    • 2019-10-11
    相关资源
    最近更新 更多