Databricks遍历Blob存储目录加载文件到Delta表及TXT Schema定义问题
解决方案
前置步骤:挂载Blob存储(如已挂载可跳过)
首先将Azure Blob存储容器挂载到Databricks文件系统,统一路径访问入口,示例代码:
# 配置Blob存储认证信息 storage_account_name = "你的存储账户名" container_name = "你的容器名" storage_account_access_key = "你的存储账户访问密钥" mount_point = f"/mnt/{container_name}" # 执行挂载操作 dbutils.fs.mount( source = f"wasbs://{container_name}@{storage_account_name}.blob.core.windows.net/", mount_point = mount_point, extra_configs = {f"fs.azure.account.key.{storage_account_name}.blob.core.windows.net": storage_account_access_key} )
步骤1:定义统一Schema
你可以提前为无预定义Schema的txt文件定义统一StructType结构,parquet文件本身自带Schema,读取时会自动和你定义的Schema做字段对齐:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType, DoubleType # 按需替换为你实际的字段名、字段类型、非空约束配置 custom_schema = StructType([ StructField("id", StringType(), nullable=True), StructField("event_time", TimestampType(), nullable=True), StructField("metric_value", DoubleType(), nullable=True), StructField("source_tag", StringType(), nullable=True) ])
如果无法确定txt文件的字段结构,可以先采样少量txt文件推断Schema,确认后再固化为上面的自定义Schema:
# 采样单个txt文件推断Schema,确认后再做正式定义 sample_txt_path = "/mnt/你的容器名/folder/2021-10-17/14-10-20/样例文件.txt" sample_df = spark.read.option("header", "true").option("inferSchema", "true").csv(sample_txt_path, sep="\t") # 分隔符按实际情况替换 sample_df.printSchema()
步骤2:递归读取全目录下所有txt与parquet文件
开启递归遍历选项,匹配两种后缀的文件,统一用自定义Schema加载:
base_path = "/mnt/你的容器名/folder/" # 读取全目录下所有txt文件 txt_df = spark.read \ .option("recursiveFileLookup", "true") # 开启递归遍历任意层级子目录 .option("header", "true") # 若txt无表头则设为false,字段会按顺序匹配自定义Schema .option("sep", "\t") # 按需替换为txt实际分隔符,逗号分隔填"," .schema(custom_schema) \ .csv(f"{base_path}/**/*.txt") # **通配符匹配所有层级子目录 # 读取全目录下所有parquet文件,自带Schema会自动和自定义Schema对齐 parquet_df = spark.read \ .option("recursiveFileLookup", "true") \ .schema(custom_schema) \ .parquet(f"{base_path}/**/*.parquet") # 合并两类文件的DataFrame all_data_df = txt_df.unionByName(parquet_df)
步骤3:写入Delta表
根据需求选择全量覆盖或者追加写入模式:
# 全量覆盖写入Delta表 all_data_df.write.format("delta") \ .mode("overwrite") \ .save("/delta/你自定义的表存储路径") # 若需要创建表元数据方便后续SQL查询 spark.sql("CREATE TABLE IF NOT EXISTS 你的表名 USING DELTA LOCATION '/delta/你自定义的表存储路径'")
补充说明
- 如果存在txt文件表头、分隔符不统一的情况,可以先用
binaryFile接口遍历所有文件路径,按路径分组批量处理不同格式的txt文件 - 若需要保留文件来源、日期、时间维度信息,可以读取时新增
input_file_name()列获取文件完整路径,再从路径中提取日期、时间字段存入Delta表 - 数据量较大时可以开启分区写入,按日期字段分区能大幅提升后续查询效率
内容的提问来源于stack exchange,提问作者Venkatesh
相关产品推荐
相关产品推荐

