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

异步Python调用get_topic_sender向Azure Service Bus发消息过慢的优化求助

问题根源

原代码的性能瓶颈在于每次发送消息都重复创建ServiceBusClient和Topic Sender。这两个属于重量级资源,初始化时需要建立连接、握手、权限验证等操作,重复执行会产生大量额外开销——你测试的100条消息耗时120秒,绝大部分时间都消耗在对象初始化上,而非消息发送本身。

优化核心

Azure Service Bus Python SDK中的ServiceBusClient和Topic Sender(无论同步还是异步版本)都是协程安全、可长期重用的对象。只需初始化一次,后续所有消息发送复用同一个实例,就能把性能拉到和同步方案一致的水平。

基于Dependency Injector的异步实现方案

以下是通过Dependency Injector注入可重用异步Topic Sender的完整方案:

1. 安装依赖

pip install azure-servicebus dependency-injector

2. 配置DI容器与资源重用

创建容器类,将ServiceBusClient和Topic Sender注册为单例,确保整个应用生命周期内只初始化一次:

from dependency_injector import containers, providers
from azure.servicebus.aio import ServiceBusClient, ServiceBusSender

CONNECTION_STR = "你的Service Bus连接字符串"
TOPIC_NAME = "目标Topic名称"

class Container(containers.DeclarativeContainer):
    # 单例模式注册ServiceBusClient
    servicebus_client = providers.Singleton(
        ServiceBusClient.from_connection_string,
        conn_str=CONNECTION_STR,
        logging_enable=True
    )

    # 单例模式注册Topic Sender,依赖已创建的ServiceBusClient
    topic_sender = providers.Singleton(
        lambda client: client.get_topic_sender(topic_name=TOPIC_NAME),
        client=servicebus_client
    )

3. 编写消息发送服务

创建服务类,注入Topic Sender,处理逐条发送消息的逻辑:

class MessageSenderService:
    def __init__(self, sender: ServiceBusSender):
        self.sender = sender

    async def send_single_message(self, message_content: str):
        from azure.servicebus import ServiceBusMessage
        message = ServiceBusMessage(message_content)
        await self.sender.send_messages(message)

    async def send_from_iterator(self, message_iterator):
        # 处理来自迭代器的逐条消息
        for msg_content in message_iterator:
            await self.send_single_message(msg_content)

4. 使用容器执行发送

初始化容器,获取服务实例,处理迭代器中的消息:

async def main():
    container = Container()
    container.init_resources()  # 初始化所有单例资源

    sender_service = container.topic_sender()
    service = MessageSenderService(sender_service)

    # 模拟消息迭代器
    def message_iterator():
        for i in range(100):
            yield f"测试消息 {i}"

    await service.send_from_iterator(message_iterator())

    # 关闭资源
    await container.servicebus_client().close()

if __name__ == "__main__":
    import asyncio
    asyncio.run(main())
关键注意事项
  • 资源释放:应用退出时必须调用ServiceBusClient.close()释放连接资源,避免泄漏。
  • 单例安全:ServiceBusClient和Sender的单例注册在异步环境下是安全的,SDK内部已做协程安全处理。
  • 迭代器处理:即使消息来自迭代器,只要复用同一个Sender,每条消息的发送开销仅为网络传输和Service Bus处理时间,不会再包含初始化开销。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 16:25:15