如何解决从Azure容器下载Blob时存在性检查耗时过长问题?
问题描述
我有一个Azure容器,其中存储了数千个Blob,每个Blob的路径格式为id/year/month/day/hour_minute_second/file.json。我需要使用Python下载指定id、且时间在start_date到end_date范围内的所有file.json文件。目前我使用azure Python包中的BlobServiceClient实现该功能,在下载每个JSON文件前,会通过get_blob_client(blob=blob_dir).exists()方法检查Blob路径是否存在,但这个存在性检查导致整体耗时过长。
现有实现代码:
from azure.storage.blob import BlobServiceClient import pandas as pd from itertools import product class AzureContainerClient(object): def __init__(self, account_url="https://mystorage.blob.core.windows.net/", container_name='json', credential="azure"): self.account_url = account_url self.container_name = container_name self.credential = credential # Connect to the Azure self.__connect() def __connect(self): """ Connect to Azure container_name. :return: """ self.blob_client_server = BlobServiceClient(account_url=self.account_url, credential=self.credential) self.container_client = self.blob_client_server.get_container_client(container=self.container_name) def close(self) -> None: """ Close the connection. :return: None """ self.blob_client_server.close() def is_exist(self, blob: str) -> bool: """ Return True if blob exist in self.container_name else False :param blob: blob address :return: """ return self.container_client.get_blob_client(blob=blob).exists() def read_blob(self, blob) -> dict: """ Read the blob from container_client. :param blob: blob directory :return: """ data = self.container_client.get_blob_client(blob=blob).download_blob().readall() # Load the binary data into json # data = json.loads(data) return data def get_files(ids: list, start_date: str, end_date: str) -> pd.DataFrame: """ Get the json files for ids from start_date to end_date. :param ids: :param start_date: :param end_date: :return: """ date_range = pd.date_range(start=start_date, end=end_date, freq='H') # Get the generator directories for each id between start_date and end_date stores_dates_gen = product(ids, date_range) azure_container_client = AzureContainerClient() data_list = [] for id_date in stores_dates_gen: # Get the blob directory id_date_blob = f'{id_date[0]}/{"/".join(id_date[1].strftime("%Y-%m-%d-%H_%M_%S").split("-"))}/file.json' # Check the existence of the id_date blob in container if not azure_container_client.is_exist(id_date_blob): continue data = azure_container_client.read_blob(blob=id_date_blob) data_list.append((id_date[0], id_date[1], data)) df = pd.DataFrame(data=data_list, columns=['id', 'dateTime', 'data']) return df
优化方案
核心思路是避免逐个检查Blob存在性,改为批量枚举符合条件的Blob或直接尝试下载并捕获异常,减少API调用次数。
方案1:批量枚举+过滤(最优)
利用Azure Blob存储的前缀过滤功能,先枚举指定id下的所有Blob,再解析路径中的时间信息筛选出符合日期范围的文件,直接下载,完全跳过存在性检查。
优化后的代码
from azure.storage.blob import BlobServiceClient import pandas as pd from datetime import datetime class AzureContainerClient(object): def __init__(self, account_url="https://mystorage.blob.core.windows.net/", container_name='json', credential="azure"): self.account_url = account_url self.container_name = container_name self.credential = credential self.__connect() def __connect(self): self.blob_client_server = BlobServiceClient(account_url=self.account_url, credential=self.credential) self.container_client = self.blob_client_server.get_container_client(container=self.container_name) def close(self) -> None: self.blob_client_server.close() def read_blob(self, blob) -> bytes: return self.container_client.get_blob_client(blob=blob).download_blob().readall() def list_filtered_blobs(self, target_ids: list, start_dt: datetime, end_dt: datetime): """批量枚举符合条件的Blob""" filtered_blobs = [] for target_id in target_ids: # 前缀过滤:只枚举当前id下的所有Blob prefix = f"{target_id}/" for blob in self.container_client.list_blobs(name_starts_with=prefix): # 解析Blob路径中的时间部分 path_parts = blob.name.split('/') if len(path_parts) < 6 or path_parts[-1] != 'file.json': continue # 拼接时间字符串并转换为datetime对象 date_part = f"{path_parts[1]}-{path_parts[2]}-{path_parts[3]} {path_parts[4].replace('_', ':')}" try: blob_dt = datetime.strptime(date_part, "%Y-%m-%d %H:%M:%S") except ValueError: continue # 跳过格式异常的Blob # 检查时间是否在目标范围内 if start_dt <= blob_dt <= end_dt: filtered_blobs.append((target_id, blob_dt, blob.name)) return filtered_blobs def get_files(ids: list, start_date: str, end_date: str) -> pd.DataFrame: # 请根据实际传入的日期格式调整strptime的格式字符串 start_dt = datetime.strptime(start_date, "%Y-%m-%d %H:%M:%S") end_dt = datetime.strptime(end_date, "%Y-%m-%d %H:%M:%S") azure_container_client = AzureContainerClient() try: # 获取所有符合条件的Blob filtered_blobs = azure_container_client.list_filtered_blobs(ids, start_dt, end_dt) data_list = [] for target_id, blob_dt, blob_name in filtered_blobs: data = azure_container_client.read_blob(blob_name) data_list.append((target_id, blob_dt, data)) df = pd.DataFrame(data=data_list, columns=['id', 'dateTime', 'data']) return df finally: azure_container_client.close()
优势
- 仅对每个
id发起一次批量枚举请求,替代原方案中每个时间点的存在性检查请求,大幅减少API调用次数 - 直接从枚举结果中筛选符合时间范围的Blob,避免无效的存在性验证
方案2:捕获下载异常(极简改动)
如果不想修改枚举逻辑,可以去掉is_exist检查,直接尝试下载,捕获ResourceNotFoundError异常,跳过不存在的Blob。这种方式将每个Blob的API调用从2次减少到1次。
优化后的get_files函数
from azure.storage.blob.exceptions import ResourceNotFoundError from itertools import product def get_files(ids: list, start_date: str, end_date: str) -> pd.DataFrame: date_range = pd.date_range(start=start_date, end=end_date, freq='H') stores_dates_gen = product(ids, date_range) azure_container_client = AzureContainerClient() data_list = [] for id_date in stores_dates_gen: # 优化日期格式拼接,避免split和join操作 id_date_blob = f'{id_date[0]}/{id_date[1].strftime("%Y/%m/%d/%H_%M_%S")}/file.json' try: data = azure_container_client.read_blob(blob=id_date_blob) data_list.append((id_date[0], id_date[1], data)) except ResourceNotFoundError: continue # 跳过不存在的Blob df = pd.DataFrame(data=data_list, columns=['id', 'dateTime', 'data']) return df
优势
- 代码改动极小,快速见效
- 每个Blob仅需一次API调用(下载),相比原方案减少50%的请求量
额外优化建议
- 并行下载:使用
concurrent.futures.ThreadPoolExecutor实现多线程并行下载,Azure SDK本身线程安全,可进一步提升速度 - 格式简化:原代码中
strftime("%Y-%m-%d-%H_%M_%S").split("-")可直接改为strftime("%Y/%m/%d/%H_%M_%S"),避免字符串拆分拼接的额外开销
内容的提问来源于stack exchange,提问作者Mohammadreza Riahi
相关产品推荐
相关产品推荐

