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

如何使用PyArrow 15.0.0预索引Azure存储中的Parquet数据集?

优化Azure Parquet数据集查询性能的预索引方案

方案1:PyArrow新Dataset API + 元数据缓存

PyArrow的新一代Dataset API支持预生成并复用分区元数据,避免每次初始化都遍历所有文件:

import adlfs
from azure.identity import DefaultAzureCredential
import pyarrow.dataset as ds
import pickle

# 初始化Azure存储文件系统
credentials = DefaultAzureCredential()
fs = adlfs.AzureBlobFileSystem(account_name=account_name, credential=credentials)

# 元数据缓存路径(可存本地或Azure存储)
cache_path = "parquet_metadata_cache.pkl"

try:
    # 加载预生成的元数据缓存
    with open(cache_path, "rb") as f:
        dataset_meta = pickle.load(f)
    dataset = ds.dataset(
        container_name,
        format="parquet",
        filesystem=fs,
        partitioning="hive",
        metadata=dataset_meta
    )
except FileNotFoundError:
    # 首次扫描数据集,收集分区与文件统计量
    dataset = ds.dataset(
        container_name,
        format="parquet",
        filesystem=fs,
        partitioning="hive",
        collect_statistics=True  # 生成文件级统计,加速过滤
    )
    # 保存元数据缓存供后续使用
    dataset_meta = dataset.metadata
    with open(cache_path, "wb") as f:
        pickle.dump(dataset_meta, f)

# 执行过滤与列修剪查询
df = dataset.to_table(
    filter=(ds.field("exp").isin(["ZY9876", "AB1234"]) & 
            ds.field("speed") == "50" & 
            ds.field("temperature") > 30),
    columns=["需要的列1", "需要的列2"]
).to_pandas()
  • 核心优势:基于PyArrow原生新API,支持列过滤和分区过滤,缓存一次后永久复用,彻底消除首次遍历开销。
  • 注意:数据集若有更新,需重新生成缓存;可将缓存文件上传至Azure存储,实现多客户端共享。

方案2:Delta Lake元数据索引

Delta Lake会自动维护Parquet文件的元数据与索引,静态数据集转换为Delta格式后,查询性能会大幅提升:

from deltalake import DeltaTable
import adlfs

# 初始化Azure存储文件系统
fs = adlfs.AzureBlobFileSystem(account_name=account_name, credential=credentials)

# 将现有Parquet数据集转换为Delta格式(仅需执行一次)
DeltaTable.convert_to_delta(
    f"az://{container_name}",
    partition_cols=["exp", "date", "speed", "load"],
    filesystem=fs
)

# 后续查询直接使用DeltaTable
dt = DeltaTable(f"az://{container_name}", filesystem=fs)
df = dt.to_pandas(
    filters=[
        ("exp", "in", ["ZY9876", "AB1234"]),
        ("speed", "==", "50"),
        ("temperature", ">", 30)
    ],
    columns=["需要的列1", "需要的列2"]
)
  • 核心优势:自动维护元数据,支持快速分区过滤和列修剪,无需手动管理缓存;即使后续有少量数据新增,也能高效更新索引。
  • 注意:转换过程需要一次全量扫描,但后续查询性能有质的提升;需安装deltalake包。

方案3:手动构建文件索引(轻量方案)

如果不想引入额外依赖,可以手动生成文件索引,查询时先过滤索引再加载目标文件:

import adlfs
import pandas as pd
import pyarrow.parquet as pq

# 初始化Azure存储文件系统
credentials = DefaultAzureCredential()
fs = adlfs.AzureBlobFileSystem(account_name=account_name, credential=credentials)

# 文件索引路径
index_path = "parquet_file_index.csv"

try:
    file_index = pd.read_csv(index_path)
except FileNotFoundError:
    # 遍历所有Parquet文件,解析Hive分区信息
    file_paths = fs.glob(f"{container_name}/**/*.parquet")
    file_index_list = []
    for path in file_paths:
        partition_info = {}
        for segment in path.split("/"):
            if "=" in segment:
                key, val = segment.split("=", 1)
                partition_info[key] = val
        partition_info["file_path"] = path
        file_index_list.append(partition_info)
    file_index = pd.DataFrame(file_index_list)
    file_index.to_csv(index_path, index=False)

# 过滤索引,得到需要加载的文件路径
target_files = file_index[
    (file_index["exp"].isin(["ZY9876", "AB1234"])) &
    (file_index["speed"] == "50")
]["file_path"].tolist()

# 加载目标文件并应用行过滤
result_tables = []
for file in target_files:
    table = pq.read_table(file, filesystem=fs, columns=["需要的列1", "需要的列2"])
    table = table.filter(table["temperature"] > 30)
    result_tables.append(table)
df = pd.concat([t.to_pandas() for t in result_tables], ignore_index=True)
  • 核心优势:轻量无额外依赖,完全可控;适合完全静态、无需频繁更新的数据集。
  • 注意:需手动解析分区路径,文件内过滤需逐个处理,性能略逊于前两种方案;数据集更新时需手动重建索引。

内容的提问来源于stack exchange,提问作者jensr

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 00:56:26