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

如何在Azure Databricks中遍历Blob存储路径并将Avro文件转为DataFrame

批量迁移Event Hubs生成的Avro文件到合理存储位置

在Azure Databricks中完全可以实现批量遍历、读取并迁移这类分层存储的Avro文件,无需逐个处理单个文件,以下是具体实现方案:

1. 批量读取所有Avro文件

Spark原生支持递归读取目录下的所有符合格式的文件,你可以直接指定根目录或使用通配符匹配所有文件:

# 方式1:直接读取根目录下所有子目录的Avro文件(Spark会自动递归扫描)
df = spark.read.format("com.databricks.spark.avro") \
    .option("mergeSchema", "true")  # 若存在schema演变,开启自动合并schema
    .load("/mnt/mount-name/eventhub/eventhubservice/")

# 方式2:使用通配符精准匹配路径结构
df = spark.read.format("com.databricks.spark.avro") \
    .option("mergeSchema", "true")
    .load("/mnt/mount-name/eventhub/eventhubservice/*/*/*/*/*/*/*.avro")

2. 提取路径中的元数据(可选但推荐)

原路径中的PartitionId、日期等信息是重要的业务元数据,可以提取到DataFrame中,方便后续按合理维度分区存储:

from pyspark.sql.functions import input_file_name, split, element_at, regexp_extract

# 从文件路径中提取PartitionId、Year、Month、Day等字段
df_with_meta = df.withColumn("file_path", input_file_name()) \
    .withColumn("path_parts", split("file_path", "/")) \
    .withColumn("PartitionId", element_at("path_parts", 6))  # 对应示例路径中的"0"
    .withColumn("Year", element_at("path_parts", 7))        # 对应"2022"
    .withColumn("Month", element_at("path_parts", 8))       # 对应"01"
    .withColumn("Day", element_at("path_parts", 9))         # 对应"01"
    .withColumn("Hour", element_at("path_parts", 10))       # 对应"01"
    # 提取秒数(去掉.avro后缀)
    .withColumn("Second", regexp_extract(element_at("path_parts", 12), r"(\d+)\.avro", 1)) \
    .drop("file_path", "path_parts")

注意:element_at的索引需要根据你的实际路径拆分结果调整,可先通过display(df.select(input_file_name()))查看路径结构后确认索引位置。

3. 迁移到合理存储位置

推荐使用Delta Lake(Databricks原生支持)或Parquet作为目标存储格式,并按Hive风格分区(如Year/Month/Day/PartitionId)存储,提升后续查询效率:

方案A:写入Delta Lake(推荐)

df_with_meta.write \
    .mode("append")  # 数据量大时用append,首次迁移可改用overwrite
    .partitionBy("Year", "Month", "Day", "PartitionId")  # 按合理维度分区
    .format("delta")
    .save("/mnt/mount-name/new-storage/eventhub-delta/")

方案B:写入Parquet

df_with_meta.write \
    .mode("append")
    .partitionBy("Year", "Month", "Day", "PartitionId")
    .format("parquet")
    .save("/mnt/mount-name/new-storage/eventhub-parquet/")

4. 超大数据集的分批处理(可选)

如果数据量极大,一次性读取内存压力大,可以按目录层级分批遍历处理:

base_path = "/mnt/mount-name/eventhub/eventhubservice/"
# 遍历所有分区目录(0和1)
partitions = [p.path for p in dbutils.fs.ls(base_path) if p.isDir]

for partition_dir in partitions:
    # 遍历年份目录
    years = [y.path for y in dbutils.fs.ls(partition_dir) if y.isDir]
    for year_dir in years:
        # 遍历月份目录
        months = [m.path for m in dbutils.fs.ls(year_dir) if m.isDir]
        for month_dir in months:
            # 遍历日期目录
            days = [d.path for d in dbutils.fs.ls(month_dir) if d.isDir]
            for day_dir in days:
                # 读取当天所有Avro文件
                daily_df = spark.read.format("com.databricks.spark.avro") \
                    .option("mergeSchema", "true") \
                    .load(f"{day_dir}/*/*/*.avro")
                
                # 提取元数据(同步骤2)
                daily_df_with_meta = daily_df.withColumn("file_path", input_file_name()) \
                    .withColumn("path_parts", split("file_path", "/")) \
                    .withColumn("PartitionId", element_at("path_parts", 6)) \
                    .withColumn("Year", element_at("path_parts", 7)) \
                    .withColumn("Month", element_at("path_parts", 8)) \
                    .withColumn("Day", element_at("path_parts", 9)) \
                    .drop("file_path", "path_parts")
                
                # 追加写入目标存储
                daily_df_with_meta.write \
                    .mode("append") \
                    .partitionBy("Year", "Month", "Day", "PartitionId") \
                    .format("delta") \
                    .save("/mnt/mount-name/new-storage/eventhub-delta/")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 05:55:18