如何修改Polars与adlfs代码,实现向Azure存储写入Delta目录?
问题
我们可以通过polars直接将parquet文件写入Azure存储(如基础存储容器)。但因业务需求需写入基于parquet的Delta格式,参考Polars官方文档对Delta写入的支持,修改代码后出现问题:
- 当指定
path = path/to/1.parquet时,可正常写入Azure存储; - 当指定
path = path/to/delta_folder/时,仅在Azure存储的delta_folder中生成一个0字节文件; - 直接在本地文件系统调用
pdf.write_delta(path, mode="append")可正常工作。
原代码如下:
import adlfs import polars as pl from azure.identity.aio import DefaultAzureCredential # pdf: pl.DataFrame # path: str # account_name: str # container_name: str credential = DefaultAzureCredential() fs = adlfs.AzureBlobFileSystem(account_name=account_name, credential=credential) with fs.open(f"{container_name}/way/to/{path}", mode="wb") as f: if path.endswith(".parquet"): pdf.write_parquet(f) else: pdf.write_delta(f, mode="append")
请问如何修改代码以支持向云端的delta_folder/进行递归写入?
解决方案
问题核心在于:pl.DataFrame.write_delta()并不支持传入单个文件对象(如fs.open()返回的句柄),因为Delta Lake格式本身是目录级存储,需要创建多个文件(数据文件、事务日志等),而非单个文件。而write_parquet()支持单文件写入,所以之前的写法对Parquet有效,但对Delta无效。
正确的做法是直接让Polars使用adlfs提供的文件系统来操作Azure Blob存储,无需手动打开文件句柄:
修改后的代码
import adlfs import polars as pl from azure.identity.aio import DefaultAzureCredential # pdf: pl.DataFrame # path: str # account_name: str # container_name: str credential = DefaultAzureCredential() # 初始化Azure Blob文件系统 fs = adlfs.AzureBlobFileSystem(account_name=account_name, credential=credential) # 构造完整的云端路径:abfs://容器名/路径 full_path = f"abfs://{container_name}/way/to/{path}" if path.endswith(".parquet"): pdf.write_parquet(full_path, filesystem=fs) else: pdf.write_delta(full_path, mode="append", filesystem=fs)
关键说明
- 使用abfs协议路径:Polars支持直接识别
abfs://开头的Azure Blob存储路径,配合传入filesystem参数指定adlfs.AzureBlobFileSystem实例,就能直接操作云端目录。 - Delta格式的目录特性:
write_delta()会自动在指定的delta_folder/目录下生成所需的Parquet数据文件、_delta_log/事务日志目录及相关文件,完整符合Delta Lake的存储规范。 - 权限与异步处理:如果使用
DefaultAzureCredential的异步版本,确保代码运行在异步上下文(如asyncio.run()包裹)中;若无需异步,可改用同步的DefaultAzureCredential(从azure.identity导入,而非azure.identity.aio)。
内容的提问来源于stack exchange,提问作者Michel Hua
相关产品推荐
相关产品推荐

