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

RabbitMQ宕机后Python基于aioamqp的重连问题排查

解决aioamqp连接RabbitMQ宕机后挂起无法重连的问题

问题现状

使用aioamqp编写的Python RPC消费者,在RabbitMQ意外宕机后,代码会挂起在连接关闭环节,无法自动触发重连逻辑,只能手动重启服务。日志显示连接被重置后,没有进入用户代码的异常处理或重连流程。

核心原因

  1. 等待逻辑未被打断:代码在await on_close.wait()处阻塞,当RabbitMQ宕机触发ConnectionResetError时,aioamqp内部的连接丢失事件不会主动打断这个await,导致代码一直卡在等待状态,无法进入重连循环。
  2. 关闭操作无超时保护:当连接已失效时,await protocol.close()可能因等待服务器响应而无限挂起,无法完成关闭操作。
  3. 异常捕获覆盖不足:aioamqp在连接丢失时抛出的部分异常未被当前except块捕获,导致无法触发重连。

解决方案

修改后的RPC服务器代码

async def rpc_server(url: str, exchange: str, queue_name: str, service_id: str, on_close: asyncio.Event):
    log.info('amqp: %s, exchange: %s, queue_name: %s, service_id: %s', url, exchange, queue_name, service_id)
    if '%2f' in url:
        url = url.replace('%2f', '/')

    o = urlparse(url)    

    backoff = 1.0
    max_attempts = 5
    retry_interval = 1

    while True:
        connection_lost = asyncio.Event()
        transport = None
        protocol = None
        try:
            transport, protocol = await aioamqp.connect(
                host=o.hostname, 
                port=o.port, 
                login=o.username, 
                password=o.password,
                virtualhost=o.path[1:]
            )

            # 监听连接意外关闭事件
            def on_connection_lost(exc):
                log.error('AMQP connection lost unexpectedly: %s', exc)
                connection_lost.set()

            protocol.add_close_callback(on_connection_lost)

            channel = await protocol.channel()
            await channel.queue_declare(queue_name=queue_name)
            await channel.queue_bind(queue_name=queue_name, exchange_name=exchange, routing_key=queue_name)
            await channel.queue_bind(queue_name=queue_name, exchange_name=exchange, routing_key=service_id)
            await channel.basic_consume(on_request, queue_name=queue_name)
            log.info('Awaiting RPC requests')

            # 同时等待正常退出事件和意外断开事件
            done, pending = await asyncio.wait(
                [on_close.wait(), connection_lost.wait()],
                return_when=asyncio.FIRST_COMPLETED
            )

            # 处理正常退出逻辑
            if on_close.is_set():
                log.info('Initiating normal AMQP connection close')
                try:
                    # 给关闭操作添加超时,避免挂起
                    await asyncio.wait_for(protocol.close(), timeout=5)
                except asyncio.TimeoutError:
                    log.error('Timeout closing AMQP protocol, force closing transport')
                finally:
                    if transport and not transport.is_closing():
                        transport.close()
                break
            # 处理意外断开逻辑,直接进入重连
            else:
                log.info('Connection lost, preparing to reconnect')
                if transport and not transport.is_closing():
                    transport.close()

        except (OSError, aioamqp.exceptions.AmqpClosedConnection, ConnectionResetError) as e:
            log.error('AMQP connection error: %s', e)
            if not max_attempts:
                raise ProcessorException('AMQP error: max attempt reached: %s', e)

            log.info('Retry in %s seconds', backoff)        
            await asyncio.sleep(backoff)
            max_attempts -= 1
            backoff *= 2
        except Exception as e:
            log.error('Unexpected exception: %s', e)
            import traceback
            traceback.print_exc()
            if max_attempts > 0:
                log.info('Retry in %s seconds', backoff)        
                await asyncio.sleep(backoff)
                max_attempts -= 1
                backoff *= 2
            else:
                raise ProcessorException('Unexpected error after max retries: %s', e)
        else:
            log.info("Normal exit from connection loop")
            break
    log.info("Exited RPC server loop")

修改后的主函数代码

def main(*args, **kwargs):
    log.info('main args: %s', kwargs)
    
    close_event = asyncio.Event()
    event_loop = asyncio.get_event_loop()
    try:
        event_loop.create_task(sim_server.run())
        event_loop.run_until_complete(rpc_server(**kwargs, on_close=close_event))
        log.info("After event_loop.run_until_complete")
    except ProcessorException as e:
        log.error('Processor exception: %s', e)
    except Exception as e:
        log.exception('Unexpected exception: %s', e)
    except KeyboardInterrupt:
        log.info("Received KeyboardInterrupt, shutting down")
        close_event.set()
    finally:
        log.info("Inside finally")
        close_event.set()
        # 等待所有任务完成后关闭事件循环
        pending = asyncio.all_tasks(event_loop)
        if pending:
            event_loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True))
        event_loop.close()

关键修改点说明

  1. 添加连接丢失回调:通过protocol.add_close_callback()监听连接关闭事件,当连接意外断开时设置connection_lost事件,打断阻塞的等待逻辑。
  2. 双事件等待机制:使用asyncio.wait()同时监听正常退出的on_close和意外断开的connection_lost事件,优先响应先触发的事件。
  3. 关闭操作超时保护:用asyncio.wait_for()给protocol.close()添加5秒超时,避免因连接失效导致的无限挂起。
  4. 优化异常处理:合并同类异常捕获,同时在意外异常时执行重试逻辑(直到达到最大次数)。
  5. 修复主函数阻塞问题:将event_loop.run_forever()替换为等待所有pending任务完成,避免无限阻塞。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 11:53:19