如何通过Databricks(PySpark)直接向Azure DataLake写入二进制文件?
在Databricks中使用PySpark将二进制docx文件写入Azure Data Lake Storage的解决方案
错误原因分析
- 使用
open()直接访问ADLS路径:open()仅支持本地文件系统操作,无法识别ADLS的adl://或abfss://协议路径,因此抛出FileNotFoundError。 dbutils.fs.put()传入二进制数据:该方法默认接受字符串类型参数,直接传入bytes对象会触发类型错误。- base64编码写入:保存的是编码后的字符串而非原始二进制字节流,破坏了docx文件的二进制结构,导致文件无法正常读取。
可行解决方案
方案1:使用Spark Binary File格式写入
Spark提供的binaryFile格式可直接处理二进制文件写入,适合通过PySpark批量或单文件写入场景:
from pyspark.sql import Row from pyspark.sql.types import StructType, StructField, BinaryType, StringType # 从Salesforce获取的docx二进制内容 binary_data = request.content # ADLS路径(Gen1用adl://,Gen2推荐用abfss://) target_path = "abfss://<container>@<storage-account>.dfs.core.windows.net/<folders>/Report.docx" # 构建包含文件路径和二进制内容的DataFrame schema = StructType([ StructField("path", StringType(), nullable=False), StructField("content", BinaryType(), nullable=False) ]) df = spark.createDataFrame([Row(path=target_path, content=binary_data)], schema=schema) # 写入二进制文件 df.write.format("binaryFile") \ .option("path", target_path) \ .mode("overwrite") \ .save()
方案2:使用Azure Data Lake Storage SDK直接写入
通过官方SDK直接操作ADLS,适合对写入过程有更精细控制的场景:
- 先安装SDK(若集群未预装):
%pip install azure-storage-file-datalake
- 写入代码:
from azure.storage.filedatalake import DataLakeServiceClient # 初始化客户端(支持账号密钥或服务主体认证) account_name = "<你的存储账户名>" account_key = "<你的存储账户密钥>" service_client = DataLakeServiceClient( account_url=f"https://{account_name}.dfs.core.windows.net", credential=account_key ) # 获取文件系统、目录客户端 file_system_client = service_client.get_file_system_client(file_system="<你的容器名>") directory_client = file_system_client.get_directory_client("<目标文件夹路径>") # 创建文件并写入二进制数据 file_client = directory_client.create_file("Report.docx") file_client.upload_data(binary_data, overwrite=True)
方案3:临时文件中转(适用于小文件)
先将二进制数据写入Databricks本地临时文件,再复制到ADLS:
import tempfile import os binary_data = request.content target_path = "adl://<something>.azuredatalakestore.net/<...folders...>/Report.docx" # 写入本地临时文件 with tempfile.NamedTemporaryFile(mode="wb", delete=False) as tmp_file: tmp_file.write(binary_data) tmp_local_path = tmp_file.name # 复制到ADLS dbutils.fs.cp(f"file:{tmp_local_path}", target_path, overwrite=True) # 清理本地临时文件 os.unlink(tmp_local_path)
内容的提问来源于stack exchange,提问作者Debtanu Gupta
相关产品推荐
相关产品推荐

