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

如何使用Python监听Microsoft.Storage.BlobCreated事件

Python监控Azure Blob Storage Blob创建事件实现方案

你需要监听的Microsoft.Storage.BlobCreated事件,和Logic Apps里的Blob创建触发逻辑完全对齐,有两种成熟实现方式,可根据场景选择。

参考的Logic Apps触发配置样式:
Logic Apps BlobCreated触发配置


前置依赖安装

先安装需要的Azure SDK包:

pip install azure-eventgrid azure-storage-blob azure-identity azure-mgmt-eventgrid aiohttp

权限准备:

  • 给运行代码的身份(本地调试可用Azure CLI登录态,云上可用托管标识/服务主体)分配存储账号的Storage Blob Data Reader、EventGrid EventSubscription Contributor权限
  • 确保存储账号所在订阅已注册Event Grid资源提供程序

方案1:Event Grid事件推送(生产推荐,和Logic Apps触发原理完全一致)

这个方案是生产环境首选,原理和Logic Apps的触发逻辑完全相同:存储账号产生Blob创建事件后,主动推送到你的Python服务,事件延迟通常在1s以内,不需要循环请求存储接口。

实现步骤

  1. 编写HTTP服务接收Event Grid推送的事件,必须先处理Event Grid的端点校验逻辑(首次绑定事件订阅时,Event Grid会发送校验请求,不返回正确校验码订阅会直接创建失败,是高频踩坑点)
  2. 为存储账号创建事件订阅,指定只订阅Microsoft.Storage.BlobCreated事件,将接收端点配置为你的Python服务地址
  3. 解析收到的事件内容,提取Blob地址、大小等信息,执行自定义业务逻辑

完整代码示例

from aiohttp import web
from azure.identity import DefaultAzureCredential
from azure.mgmt.eventgrid import EventGridManagementClient
from azure.mgmt.eventgrid.models import (
    WebHookEventSubscriptionDestination,
    EventSubscriptionFilter
)

# 替换为你自己的配置
AZURE_SUB_ID = "你的Azure订阅ID"
STORAGE_RG = "存储账号所在资源组名称"
STORAGE_ACCOUNT = "要监控的存储账号名"
EVENT_SUB_NAME = "python-blob-created-monitor"
LISTEN_PORT = 8080
# 本地调试可用ngrok将本地端口映射为公网地址,生产环境替换为服务实际公网/内网可达地址
EVENT_RECEIVE_ENDPOINT = f"http://你的服务地址:{LISTEN_PORT}/event"

async def event_receive_handler(request):
    event_list = await request.json()
    for event in event_list:
        # 处理Event Grid端点校验请求
        if request.headers.get("aeg-event-type") == "SubscriptionValidation":
            validation_code = event["data"]["validationCode"]
            return web.json_response({"validationResponse": validation_code})
        
        # 处理Blob创建事件
        if event["eventType"] == "Microsoft.Storage.BlobCreated":
            event_data = event["data"]
            print("="*50)
            print(f"检测到新Blob创建:")
            print(f"Blob完整URL:{event_data['url']}")
            print(f"触发操作API:{event_data['api']}")
            print(f"Blob大小:{event_data['contentLength']}字节")
            print(f"内容类型:{event_data['contentType']}")
            # 在此处添加你的自定义业务逻辑,比如下载Blob、触发后续处理流程
    return web.Response(status=200)

def init_event_subscription():
    cred = DefaultAzureCredential()
    eg_client = EventGridManagementClient(cred, AZURE_SUB_ID)
    storage_scope = f"/subscriptions/{AZURE_SUB_ID}/resourceGroups/{STORAGE_RG}/providers/Microsoft.Storage/storageAccounts/{STORAGE_ACCOUNT}"
    
    # 配置事件订阅:仅接收BlobCreated事件
    sub_destination = WebHookEventSubscriptionDestination(endpoint_url=EVENT_RECEIVE_ENDPOINT)
    sub_filter = EventSubscriptionFilter(included_event_types=["Microsoft.Storage.BlobCreated"])
    
    eg_client.event_subscriptions.begin_create_or_update(
        scope=storage_scope,
        event_subscription_name=EVENT_SUB_NAME,
        destination=sub_destination,
        filter=sub_filter
    ).result()
    print("事件订阅初始化完成")

if __name__ == "__main__":
    # 首次运行执行一次初始化即可,后续运行可以注释掉这行
    init_event_subscription()
    app = web.Application()
    app.add_routes([web.post("/event", event_receive_handler)])
    web.run_app(app, port=LISTEN_PORT)

注意事项

  • 本地调试如果没有公网IP,可以用ngrok等内网穿透工具将本地端口映射为公网HTTPS地址,填入EVENT_RECEIVE_ENDPOINT即可
  • 如果不想暴露公网端点,可以将事件投递目标换成Azure队列存储/服务总线,Python端直接从队列拉取事件,逻辑完全一致,安全性更高

方案2:轮询检测(轻量无配置,适合测试/小规模场景)

如果不想配置Event Grid、公网端点,可以直接用Blob SDK轮询容器内的Blob列表,对比增量检测新创建的Blob,不需要额外配置Azure资源,适合小规模测试场景。

完整代码示例

from azure.storage.blob import BlobServiceClient
import time

# 替换为你自己的配置
STORAGE_CONN_STR = "你的存储账号连接字符串"
MONITOR_CONTAINER = "要监控的容器名称"
POLL_INTERVAL = 10 # 轮询间隔,单位秒,不建议小于5秒避免触发限流

def poll_monitor():
    blob_service = BlobServiceClient.from_connection_string(STORAGE_CONN_STR)
    container = blob_service.get_container_client(MONITOR_CONTAINER)
    
    # 初始化已存在的Blob列表
    known_blob_set = set(blob.name for blob in container.list_blobs())
    print(f"初始化完成,已记录{len(known_blob_set)}个存量Blob,开始监控...")

    while True:
        current_blob_set = set()
        for blob in container.list_blobs():
            current_blob_set.add(blob.name)
            if blob.name not in known_blob_set:
                print("="*50)
                print(f"检测到新Blob创建:{blob.name}")
                print(f"创建时间:{blob.creation_time}")
                print(f"Blob大小:{blob.size}字节")
                # 在此处添加自定义业务逻辑
        known_blob_set = current_blob_set
        time.sleep(POLL_INTERVAL)

if __name__ == "__main__":
    poll_monitor()

注意事项

  • 轮询间隔不要设置过短,否则容易触发存储账号的API请求限流
  • 单容器Blob数量超过10万时,轮询的性能和成本会明显上升,不建议生产环境使用
  • 如果需要检测Blob被覆盖更新的场景,可以额外对比Blob的last_modified时间字段

内容的提问来源于stack exchange,提问作者ezicab

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 08:45:35