如何用Python/PySpark将Blob内容写入Microsoft Fabric OneLake
迁移至OneLake的实现方案
一、OneLake连接与文件写入代码实现
OneLake兼容Azure Data Lake Storage Gen2(ADLS Gen2)API,因此可通过Azure Storage SDK的DataLakeServiceClient完成文件操作。以下是适配后的完整代码:
1. 依赖包安装
首先确保安装所需SDK:
pip install azure-storage-file-datalake azure-identity
2. 核心代码实现
from azure.storage.filedatalake import DataLakeServiceClient from azure.identity import ClientSecretCredential import base64 from pyspark.sql.functions import udf from pyspark.sql.types import BinaryType # OneLake 配置参数 tenant_id = "<你的Azure AD租户ID>" client_id = "<你的服务主体客户端ID>" client_secret = "<你的服务主体密钥>" workspace_name = "<你的Fabric工作区名称>" lakehouse_name = "<你的Lakehouse名称>" target_filesystem = "Files" # Lakehouse默认文件存储容器,可按需调整 # 构建服务主体凭据 credential = ClientSecretCredential( tenant_id=tenant_id, client_id=client_id, client_secret=client_secret ) # 初始化OneLake服务客户端 onelake_endpoint = f"https://onelake.dfs.fabric.microsoft.com/{workspace_name}/{lakehouse_name}" datalake_client = DataLakeServiceClient(account_url=onelake_endpoint, credential=credential) # 获取目标文件系统客户端 filesystem_client = datalake_client.get_file_system_client(target_filesystem) def write_file_to_onelake(data, filename): # 自动创建文件路径中的目录(若不存在) path_components = filename.split('/') if len(path_components) > 1: dir_path = '/'.join(path_components[:-1]) filesystem_client.create_directory(dir_path, exists_ok=True) # 上传文件到OneLake file_client = filesystem_client.get_file_client(filename) file_client.upload_data(data, overwrite=True) # Base64解码UDF(与原逻辑一致) def decode_base64(base64_str): return base64.b64decode(base64_str) # 注册Spark UDF decode_udf = udf(decode_base64, BinaryType())
3. 调用方式
collected_data = df_with_decoded_data.collect() # 批量写入文件到OneLake for row in collected_data: write_file_to_onelake(row['DecodedData'], row['FinalFileName'])
二、OneLake支持的凭据类型
OneLake基于Azure AD进行身份认证,支持以下常用凭据类型:
- 服务主体(Client Secret):适合自动化批量任务,通过租户ID、客户端ID和密钥完成认证,如上述示例。
- 托管标识(Managed Identity):适用于Azure托管服务(如VM、Function)中的应用,无需手动管理凭据。
- Azure AD用户凭据:使用用户名+密码认证,适合交互式场景,但受MFA限制。
- Azure CLI凭据:本地开发时可复用已登录的Azure CLI会话,无需硬编码凭据。
- SAS令牌:生成有限权限、有效期的共享访问令牌,用于临时授权访问。
内容的提问来源于stack exchange,提问作者Sreedhar
相关产品推荐
相关产品推荐

