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

从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 20:27:03