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

如何基于PyArrow Dataset高效批量迁移Azure小Parquet文件至DeltaLake

批量小Parquet文件迁移至DeltaLake的高效处理方案

我们的Azure Blob存储中有大量小Parquet文件,想要高效、低内存占用地迁移到DeltaLake。尝试用PyArrow Dataset读取所有文件并调用to_batches方法,但PyArrow不会合并多文件生成指定batch_size的批次,而是单个文件对应一个批次。向DeltaLake追加大量小文件会产生大量微文件,合并开销极高。请问有没有高效处理文件、减少与DeltaTable交互次数的方法?能不能调整分片逻辑,让返回的DataFrame包含多个文件的数据?

原代码

def read_parquet_files_from_azure(container_name, connection_string, prefix, fs_in):
    # Create a blob service client
    blob_service_client = BlobServiceClient.from_connection_string(connection_string)

    # # Get a container client
    container_client = blob_service_client.get_container_client(container_name)

    # # Get all blob names in the container with the given prefix
    blob_names = [blob.name for blob in container_client.list_blobs(name_starts_with=prefix)]

    # # Create a list of full blob URIs
    blob_uris = [f"{container_name}/{blob_name}" for blob_name in blob_names]

    # Read all parquet files into a single PyArrow dataset
    dataset = ds.dataset(
        blob_uris,
        format="parquet",
        filesystem=fs_in,
    )

    return dataset


dataset = read_parquet_files_from_azure(
    "migration-dev", os.environ["INPUT_AZURE_STORAGE_CONNECTION_STRING"], "table/", input_fs
)

i = 0
for slice in dataset.to_batches():

    log.info(f"Shape of slice: ({slice.num_rows}, {slice.num_columns})")

    df = slice.to_pandas()

    # Extract the date from the "Date" timestamp column
    # dates = df["Date"].dt.date.values
    dates = df["Date"].dt.to_period("M").astype(str).str.replace("-", "_")  # partition by year-month

    # Convert the pandas DataFrame back to a PyArrow Table
    slice = pa.table(slice).append_column("date_id", pa.array(dates))

    # write to deltalkae table w.r.t. deltalake_write_mode
    write_deltalake(
        storage_options=output_sa,
        table_or_uri=delta_table_path,
        data=slice,
        mode="append",
        partition_by=partition_keys,
        engine="rust",
    )
    i = i + 1

    if i % 10 == 0:
        # COMPACT
        log.info("Compacting")
        delta_table = DeltaTable(
            table_uri=f"""abfs://{os.environ["CONTAINER_DELTALAKE"]}/operating_report/""",
            storage_options=output_sa,
        )
        delta_table_optimizer = TableOptimizer(delta_table)
        log.info(delta_table_optimizer.compact())

    if i % 100 == 0:
        # VACCUM
        log.info("Vacuuming")
        delta_table = DeltaTable(
            table_uri=f"""abfs://{os.environ["CONTAINER_DELTALAKE"]}/operating_report/""",
            storage_options=output_sa,
        )
        log.info(delta_table.vacuum(dry_run=False))

解决方案

1. 手动累加批次,实现跨文件合并

PyArrow默认按文件拆分批次,可通过手动累加批次数据,达到目标大小后再写入,大幅减少与DeltaLake的交互次数:

target_batch_size = 100000  # 根据内存容量调整目标行数
current_batch = None
i = 0

for batch in dataset.to_batches():
    if current_batch is None:
        current_batch = batch
    else:
        # 合并当前批次与新读取的批次
        current_batch = pa.concat_tables([current_batch, batch])
    
    # 达到目标批次大小或遍历完所有文件时执行写入
    if current_batch.num_rows >= target_batch_size or batch == list(dataset.to_batches())[-1]:
        # 处理date_id字段
        df = current_batch.to_pandas()
        dates = df["Date"].dt.to_period("M").astype(str).str.replace("-", "_")
        current_batch = current_batch.append_column("date_id", pa.array(dates))
        
        # 写入DeltaLake
        write_deltalake(
            storage_options=output_sa,
            table_or_uri=delta_table_path,
            data=current_batch,
            mode="append",
            partition_by=partition_keys,
            engine="rust",
        )
        current_batch = None
        i += 1

2. 按分区预分组文件,批量读取写入

既然已按date_id分区,可先按日期分组文件,每组文件一次性读取成大批次写入,避免跨分区的无效合并:

from collections import defaultdict

def read_parquet_files_grouped_by_date(container_name, connection_string, prefix, fs_in):
    blob_service_client = BlobServiceClient.from_connection_string(connection_string)
    container_client = blob_service_client.get_container_client(container_name)
    
    # 按日期分组文件(示例:从文件路径提取日期分区)
    file_groups = defaultdict(list)
    for blob in container_client.list_blobs(name_starts_with=prefix):
        # 假设文件路径格式为table/2023_01/data.parquet,提取日期部分
        date_part = blob.name.split("/")[1]
        file_groups[date_part].append(f"{container_name}/{blob.name}")
    
    # 为每个日期组创建数据集
    grouped_datasets = {}
    for date, uris in file_groups.items():
        grouped_datasets[date] = ds.dataset(uris, format="parquet", filesystem=fs_in)
    
    return grouped_datasets

# 使用分组后的数据集批量写入
grouped_datasets = read_parquet_files_grouped_by_date(
    "migration-dev", os.environ["INPUT_AZURE_STORAGE_CONNECTION_STRING"], "table/", input_fs
)

for date, dataset in grouped_datasets.items():
    # 读取整个日期组的数据为单个表
    table = dataset.to_table()
    df = table.to_pandas()
    dates = df["Date"].dt.to_period("M").astype(str).str.replace("-", "_")
    table = table.append_column("date_id", pa.array(dates))
    
    write_deltalake(
        storage_options=output_sa,
        table_or_uri=delta_table_path,
        data=table,
        mode="append",
        partition_by=partition_keys,
        engine="rust",
    )

3. 优化Compact与Vacuum触发策略

原代码频繁触发Compact会增加额外开销,建议改为全量写入完成后一次性执行,或按文件数量/数据量阈值触发:

# 所有文件写入完成后执行一次Major Compact与Vacuum
delta_table = DeltaTable(
    table_uri=f"""abfs://{os.environ["CONTAINER_DELTALAKE"]}/operating_report/""",
    storage_options=output_sa,
)
delta_table_optimizer = TableOptimizer(delta_table)
# Major Compact会合并所有小文件
delta_table_optimizer.compact(major=True)
delta_table.vacuum(dry_run=False)

4. (可选)用PySpark替代PyArrow处理大规模数据

如果环境允许,PySpark原生支持合并小文件,且能并行处理,效率更高:

from pyspark.sql import SparkSession
from pyspark.sql.functions import date_format, regexp_replace, col

spark = SparkSession.builder \
    .appName("ParquetToDelta") \
    .config("spark.sql.parquet.mergeSchema", "true") \
    .config("spark.delta.mergeSchema", "true") \
    .getOrCreate()

# 读取Azure Blob中的Parquet文件
df = spark.read.parquet(f"abfs://migration-dev/table/")

# 添加date_id字段
df = df.withColumn("date_id", regexp_replace(date_format(col("Date"), "yyyy_MM"), "-", "_"))

# 写入DeltaLake,自动合并小文件
df.write \
    .format("delta") \
    .mode("append") \
    .partitionBy("date_id") \
    .option("mergeSchema", "true") \
    .save(delta_table_path)

# 执行优化与清理
spark.sql(f"OPTIMIZE delta.`{delta_table_path}` ZORDER BY Date")
spark.sql(f"VACUUM delta.`{delta_table_path}`")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 02:05:55