You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何解决从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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.29 05:27:01