如何基于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
相关产品推荐
相关产品推荐

