从Blob Storage自动化提取数据至Azure Databricks方案咨询
Azure Databricks 高效加载Blob嵌套Parquet数据及存储方案
一、历史数据批量加载
针对Blob中按年/月/日嵌套存储的历史Parquet数据,可通过Spark递归加载+Delta Lake批量写入实现高效迁移:
1. 全量递归加载(数据量适中场景)
直接用通配符匹配所有嵌套路径下的Parquet文件,一次性加载后写入Delta表:
# 替换为你的Blob挂载路径 mount_path = "/dbfs/mnt/blob_data" # 递归加载所有历史数据 df_history = spark.read.parquet(f"{mount_path}/**/table_name.parquet") # 从文件路径提取加载日期并添加列 from pyspark.sql.functions import input_file_name, regexp_extract df_history = df_history.withColumn( "load_date", regexp_extract(input_file_name(), r"(\d{4})/(\d{2})/(\d{2})", 1) + "-" + regexp_extract(input_file_name(), r"(\d{4})/(\d{2})/(\d{2})", 2) + "-" + regexp_extract(input_file_name(), r"(\d{4})/(\d{2})/(\d{2})", 3) ) # 写入Delta表,按load_date分区 df_history.write.mode("overwrite").partitionBy("load_date").format("delta").save("/dbfs/delta/table_name")
2. 分批加载(超大规模数据场景)
若历史数据量过大,遍历日期文件夹分批处理,避免内存压力:
mount_path = "/dbfs/mnt/blob_data" delta_table_path = "/dbfs/delta/table_name" # 获取所有日期层级的文件夹 date_folders = [f.path for f in dbutils.fs.ls(mount_path) if f.isDir()] for folder in date_folders: # 从路径提取年、月、日(根据实际路径结构调整索引) path_parts = folder.strip('/').split('/') year, month, day = path_parts[-3], path_parts[-2], path_parts[-1] # 加载当前日期文件夹下的Parquet文件 df_batch = spark.read.parquet(f"{folder}/table_name.parquet") df_batch = df_batch.withColumn("load_date", lit(f"{year}-{month}-{day}")) # 追加到Delta表 df_batch.write.mode("append").partitionBy("load_date").format("delta").save(delta_table_path)
二、每日增量自动化加载
通过Databricks Jobs实现每日自动加载当天的Blob数据:
1. 增量加载脚本
编写Python脚本,加载当日路径下的Parquet文件并追加到Delta表:
from datetime import datetime from pyspark.sql.functions import lit mount_path = "/dbfs/mnt/blob_data" delta_table_path = "/dbfs/delta/table_name" # 获取当日日期,匹配Blob的路径格式 today = datetime.today() year = str(today.year) month = str(today.month).zfill(2) day = str(today.day).zfill(2) daily_file_path = f"{mount_path}/{year}/{month}/{day}/table_name.parquet" # 检查文件是否存在,避免空加载 if dbutils.fs.exists(daily_file_path): df_daily = spark.read.parquet(daily_file_path) df_daily = df_daily.withColumn("load_date", lit(f"{year}-{month}-{day}")) # 追加到Delta表 df_daily.write.mode("append").partitionBy("load_date").format("delta").save(delta_table_path) else: print(f"Warning: No data found for {year}-{month}-{day}")
2. 配置定时调度
在Databricks控制台创建Job:
- 选择上述脚本作为任务
- 设置调度规则(如每日凌晨2点)
- 配置对应数据量的集群规格
三、Databricks内数据存储策略
1. 存储格式:Delta Lake
优先使用Delta Lake存储所有数据,核心优势:
- 支持ACID事务,避免数据写入异常
- 内置时间旅行,可回溯任意版本的数据
- 支持增量更新、合并(Merge)操作,处理重复数据更灵活
2. 分区策略
按load_date字段分区:
- 匹配数据按天生成的特性,查询时可快速过滤指定日期范围的数据
- 降低分区扫描范围,大幅提升查询性能
3. 表类型选择
- 托管表:适合仅在Databricks内使用的场景,执行
CREATE TABLE table_name USING DELTA LOCATION '/dbfs/delta/table_name',由Databricks自动管理元数据和存储生命周期 - 外部表:适合需要与其他系统共享存储的场景,执行
CREATE EXTERNAL TABLE table_name USING DELTA LOCATION '/dbfs/delta/table_name',存储位置自主管理
4. 定期优化
- OPTIMIZE:定期运行
OPTIMIZE table_name ZORDER BY (your_key_column),合并小文件并按指定列排序,提升查询速度 - VACUUM:运行
VACUUM table_name RETAIN 7 DAYS清理旧的Delta快照(根据业务需求调整保留天数),节省存储成本
内容的提问来源于stack exchange,提问作者Saswat Ray
相关产品推荐
相关产品推荐

