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

在Azure Databricks中读取挂载Blob存储的Zip文件并写入Delta表

读取挂载Blob存储中的Zip文件并写入Delta表的解决方案

方法一:使用Spark直接读取Zip中的CSV(推荐,适用于Schema一致的情况)

Spark原生支持读取Zip压缩包内的CSV文件,无需手动解压,直接通过DBFS挂载路径访问即可:

# 读取挂载点下所有Zip文件中的CSV数据
df = spark.read.csv(
    "/mnt/azureblobstorage/*.zip",  # 替换为你的实际挂载点路径
    header=True,  # 根据CSV是否包含表头调整
    inferSchema=True,  # 自动推断字段类型,也可手动指定schema
    sep=","  # 根据CSV的分隔符调整
)

# 将数据写入Delta表(支持overwrite/append等模式)
df.write.mode("append").format("delta").saveAsTable("your_database.target_table")
# 或保存到指定路径:df.write.mode("append").format("delta").save("/mnt/delta/target_table")

说明:如果Zip包内包含多个CSV文件,Spark会自动合并所有文件内容,前提是所有CSV的字段结构一致。若结构不一致,需单独处理每个Zip包。

方法二:使用Python zipfile库处理复杂解压场景

如果需要逐个处理Zip包内的文件(比如不同CSV结构不同),可以用Python的zipfile库结合DBFS的本地访问路径(/dbfs/前缀)来操作:

import zipfile
import os
import pandas as pd

# DBFS挂载点路径
mount_point = "/mnt/azureblobstorage"
# 通过/dbfs前缀直接访问DBFS文件系统
local_dbfs_path = f"/dbfs{mount_point}"

# 遍历挂载点下所有Zip文件
for root, _, files in os.walk(local_dbfs_path):
    for file in files:
        if file.endswith(".zip"):
            zip_full_path = os.path.join(root, file)
            with zipfile.ZipFile(zip_full_path, "r") as zip_ref:
                # 遍历Zip内的所有CSV文件
                for csv_file in zip_ref.namelist():
                    if csv_file.endswith(".csv"):
                        # 读取CSV内容并转为Spark DataFrame
                        with zip_ref.open(csv_file) as csv_stream:
                            pdf = pd.read_csv(csv_stream)
                            df = spark.createDataFrame(pdf)
                            # 写入Delta表
                            df.write.mode("append").format("delta").saveAsTable("your_database.target_table")

方法三:通过临时目录+Shell命令解压(兼容原有Shell逻辑)

如果必须使用unzip命令,可先将DBFS上的Zip文件复制到集群本地临时目录,解压后再读取:

import os
import shutil
from subprocess import run

# 本地临时目录
temp_unzip_dir = "/tmp/zip_extract"
os.makedirs(temp_unzip_dir, exist_ok=True)

# 复制DBFS上的Zip文件到本地临时目录
dbutils.fs.cp("/mnt/azureblobstorage/file.zip", f"file://{temp_unzip_dir}/file.zip")

# 执行解压命令
run(["unzip", f"{temp_unzip_dir}/file.zip", "-d", temp_unzip_dir], check=True)

# 读取解压后的CSV文件
df = spark.read.csv(
    f"file://{temp_unzip_dir}/*.csv",
    header=True,
    inferSchema=True
)

# 写入Delta表
df.write.mode("append").format("delta").saveAsTable("your_database.target_table")

# 清理临时文件
shutil.rmtree(temp_unzip_dir)

问题根源说明:你遇到的unzip命令报错是因为Shell命令默认访问集群本地文件系统,而/mnt/开头的路径是DBFS挂载点,无法直接被Shell命令识别。通过/dbfs/前缀访问或复制到本地临时目录即可解决。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 10:15:39