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

Python读取Service Bus主题消息存入同名Blob容器代码问题求助

问题排查与实现方案

现有代码问题梳理

  • 未保留消息原始标识:当前代码用str(uuid.uuid4()) + ".txt"作为Blob名称,完全覆盖了原始消息的对应标识,不符合数据完整性要求
  • 消息消费逻辑不完整:注释掉了receiver.complete_message(msg),会导致消息未被标记为已消费,后续会被重复读取上传
  • 本地临时文件操作冗余:不需要先将消息内容写入本地磁盘,可以直接将消息体的字节流上传到Blob,减少不必要的IO开销
  • 无异常处理逻辑:未捕获Service Bus连接异常、Blob上传失败等场景的错误,遇到异常会直接导致程序崩溃

无Azure函数依赖的完整实现方案

该方案直接使用Azure官方Python SDK实现长连接监听Service Bus消息,不需要依赖函数服务,同时满足保留原始消息标识、保障数据完整性的要求。

前置依赖安装

执行以下命令安装需要的SDK包:

pip install azure-servicebus azure-storage-blob python-dotenv

完整可运行代码

import os
import logging
from azure.servicebus import ServiceBusClient
from azure.storage.blob import BlobServiceClient
from dotenv import load_dotenv

# 加载配置,敏感信息建议存储在.env文件,不要硬编码到代码中
load_dotenv()
SERVICE_BUS_CONN_STR = os.getenv("SERVICE_BUS_CONN_STR")
# 队列模式填队列名,Topic模式填订阅名
QUEUE_OR_SUBSCRIPTION_NAME = os.getenv("QUEUE_OR_SUBSCRIPTION_NAME")
# Topic模式填对应Topic名,队列模式留空即可
TOPIC_NAME = os.getenv("TOPIC_NAME", "")
STORAGE_CONN_STR = os.getenv("STORAGE_CONN_STR")
BLOB_CONTAINER_NAME = os.getenv("BLOB_CONTAINER_NAME")

# 初始化Service Bus客户端
if TOPIC_NAME:
    servicebus_client = ServiceBusClient.from_connection_string(
        conn_str=SERVICE_BUS_CONN_STR,
        logging_enable=True
    )
    receiver = servicebus_client.get_subscription_receiver(
        topic_name=TOPIC_NAME,
        subscription_name=QUEUE_OR_SUBSCRIPTION_NAME
    )
else:
    servicebus_client = ServiceBusClient.from_connection_string(
        conn_str=SERVICE_BUS_CONN_STR,
        logging_enable=True
    )
    receiver = servicebus_client.get_queue_receiver(
        queue_name=QUEUE_OR_SUBSCRIPTION_NAME
    )

# 初始化Blob存储客户端
blob_service_client = BlobServiceClient.from_connection_string(STORAGE_CONN_STR)

if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO)
    with receiver:
        for msg in receiver:
            try:
                msg_body = msg.get_body().decode('utf-8')
                logging.info(f"收到消息,原生消息ID:{msg.message_id}")
                # 使用消息原生ID作为Blob名称,保留JSON后缀,也可自定义规则取消息自定义属性作为名称
                blob_name = f"{msg.message_id}.json"
                # 直接上传字节流,无需生成本地临时文件
                blob_client = blob_service_client.get_blob_client(
                    container=BLOB_CONTAINER_NAME,
                    blob=blob_name
                )
                blob_client.upload_blob(msg_body.encode('utf-8'), overwrite=True)
                logging.info(f"Blob上传成功,名称:{blob_name}")
                # 标记消息为已消费,避免重复读取
                receiver.complete_message(msg)
            except Exception as e:
                logging.error(f"处理消息失败:{str(e)}")
                # 处理失败的消息放回队列,可重试消费
                receiver.abandon_message(msg)

配置说明

  • 你需要在项目根目录创建.env文件,填写对应配置项即可直接运行
  • 如果需要使用消息内自定义的字段作为Blob名称,只需要修改blob_name的取值规则,从msg.application_properties中提取你自定义的名称字段即可
  • 程序会保持长连接持续监听Service Bus消息,只要运行进程不终止就会持续工作,不需要依赖Azure Functions服务

关键优化点说明

  • 保留原始标识:默认使用Service Bus消息的原生message_id作为Blob名称,完全保障数据溯源的完整性
  • 无本地文件依赖:直接将消息字节流上传到Blob,减少不必要的本地IO操作,提升运行效率
  • 异常处理完整:处理成功才确认消息,处理失败的消息自动放回队列重试,不会丢失消息
  • 同时支持队列和Topic两种模式,根据自己的Service Bus资源类型配置即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 12:24:03