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

Azure Service Bus消息自动续期失效,处理后无法完成消息

问题分析与解决方案

核心原因

你的代码存在关键问题:同步阻塞操作卡死了asyncio事件循环。

executor(topic_name, bytes, temp_dir)是耗时40分钟的同步函数,在异步上下文(async def函数)中直接调用这类同步阻塞代码会完全占用asyncio事件循环,导致AutoLockRenewer的锁续期异步任务无法被调度执行。锁续期需要定期向Service Bus发送续期请求,但只要事件循环被同步代码卡住,这些请求就发不出去,最终消息锁过期,触发MessageLockLostError。

解决方案

把同步的executor操作放到独立线程中执行,避免阻塞事件循环。可以用Python 3.9+自带的asyncio.to_thread,或者concurrent.futures.ThreadPoolExecutor实现。

修改后的代码示例

async def receive_messages_from_topic(
    servicebus_client: ServiceBusClient,
    topic_name: str,
    executor: Any,  # type: ignore
):
    async with servicebus_client:
        # 设置续期时长为3600秒(1小时),覆盖处理耗时
        renewer = AutoLockRenewer(max_lock_renewal_duration=3600)
        async with servicebus_client.get_subscription_receiver(
            topic_name=topic_name,
            subscription_name=SUBSCRIPTION_NAME,
            auto_lock_renewer=renewer,  # type: ignore
        ) as receiver:
            logger.info(f"Listening to topic: {topic_name} for subscription: {SUBSCRIPTION_NAME}")
            while True:
                try:
                    messages = await receiver.receive_messages()
                    for message in messages:
                        logger.info(
                            f"Received message from topic: {topic_name}, msg: {str(message)}"
                        )

                        with make_temp_directory() as temp_dir:
                            try:
                                bytes_str = str(message)
                                bytes_data = ast.literal_eval(bytes_str)
                                # 将同步executor放到线程执行,不阻塞事件循环
                                await asyncio.to_thread(executor, topic_name, bytes_data, temp_dir)
                                await receiver.complete_message(message)
                            except Exception as e:
                                logger.exception(
                                    f"Error executing job for topic: {topic_name} message: {message.message_id}, {e}"
                                )
                                await receiver.dead_letter_message(message)

                    if len(messages) == 0:
                        await asyncio.sleep(60)
                except Exception as e:
                    logger.exception(f"Error processing message for topic: {topic_name}, {e}")

额外注意事项

  1. 变量名冲突:原代码用bytes作为变量名,会覆盖Python内置的bytes类型,建议改为bytes_data这类名称,避免潜在问题。
  2. 续期时长验证:确保max_lock_renewal_duration设置的时长大于消息处理总耗时(这里3600秒足够覆盖40分钟)。
  3. 线程/进程池选择:如果executor是CPU密集型任务,可考虑用ProcessPoolExecutor;IO密集型任务用ThreadPoolExecutor更合适。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 15:31:16