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

