在Databricks中读取Azure Append Blob遇连接重置错误的求助
解决方案:Databricks中处理Azure Append Blob的下载与读取问题
问题核心
你遇到的ConnectionResetError主要源于直接用Python的open操作DBFS URI(dbfs:/)的兼容性问题,同时高并发下载可能触发网络限制;而本地工作目录的文件无法读取,是因为Spark默认访问DBFS,需指定本地文件协议。
方案1:修正DBFS路径的写入方式
Databricks提供了本地文件系统到DBFS的映射路径/dbfs/,替换原代码中的dbfs:/,让Python的open能直接操作:
file_dir = data.strftime("%Y%m") file_name = data.strftime("%Y%m%d") + ".csv" # 用/dbfs/替代dbfs:/,使用本地文件系统映射路径 local_path = f"/dbfs/FileStore/{file_dir}" local_file = f"{local_path}/{file_name}" blob_name = f"{file_dir}/{file_name}" # 确保目录存在 import os os.makedirs(local_path, exist_ok=True) blob_service_client_instance = BlobServiceClient(account_url=f"https://{STORAGEACCOUNTURL}.blob.core.windows.net", credential=STORAGEACCOUNTKEY) blob_client_instance = blob_service_client_instance.get_blob_client(CONTAINERNAME, blob_name) with open(local_file, "wb") as my_blob: # 降低并发数避免连接重置 blob_data = blob_client_instance.download_blob(max_concurrency=1) blob_data.readinto(my_blob)
方案2:直接用Spark读取Azure Blob(无需下载)
跳过本地下载步骤,直接通过Spark的wasbs协议读取Append Blob,更高效且避免下载错误:
file_dir = data.strftime("%Y%m") file_name = data.strftime("%Y%m%d") + ".csv" blob_path = f"wasbs://{CONTAINERNAME}@{STORAGEACCOUNTURL}.blob.core.windows.net/{file_dir}/{file_name}" # 直接读取为DataFrame df = spark.read.csv( blob_path, storage_options={ "accountKey": STORAGEACCOUNTKEY }, header=True # 根据你的CSV格式调整 ) # 后续计算后写入Cassandra df.write.format("org.apache.spark.sql.cassandra")\ .options(table="your_table", keyspace="your_keyspace")\ .mode("append")\ .save()
方案3:解决本地工作目录的文件读取问题
如果要保留临时目录方案,读取文件时需要添加file://前缀指定本地文件系统:
import os local_work_dir = os.getcwd() local_file = f"{local_work_dir}/{file_name}" # 下载文件到本地工作目录(你的原有代码) # ... # 读取时添加file://前缀 df = spark.read.csv(f"file://{local_file}", header=True)
额外优化建议
- 降低
download_blob的max_concurrency值(比如1-3),避免因高并发导致的连接重置; - 确保Databricks集群的网络能正常访问Azure Blob存储(比如集群和存储在同一区域,或配置了正确的防火墙规则)。
内容的提问来源于stack exchange,提问作者Gabriele Sciurti
相关产品推荐
相关产品推荐

