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
相关产品推荐
相关产品推荐

