RabbitMQ宕机后Python基于aioamqp的重连问题排查
解决aioamqp连接RabbitMQ宕机后挂起无法重连的问题
问题现状
使用aioamqp编写的Python RPC消费者,在RabbitMQ意外宕机后,代码会挂起在连接关闭环节,无法自动触发重连逻辑,只能手动重启服务。日志显示连接被重置后,没有进入用户代码的异常处理或重连流程。
核心原因
- 等待逻辑未被打断:代码在
await on_close.wait()处阻塞,当RabbitMQ宕机触发ConnectionResetError时,aioamqp内部的连接丢失事件不会主动打断这个await,导致代码一直卡在等待状态,无法进入重连循环。 - 关闭操作无超时保护:当连接已失效时,
await protocol.close()可能因等待服务器响应而无限挂起,无法完成关闭操作。 - 异常捕获覆盖不足: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()
关键修改点说明
- 添加连接丢失回调:通过
protocol.add_close_callback()监听连接关闭事件,当连接意外断开时设置connection_lost事件,打断阻塞的等待逻辑。 - 双事件等待机制:使用
asyncio.wait()同时监听正常退出的on_close和意外断开的connection_lost事件,优先响应先触发的事件。 - 关闭操作超时保护:用
asyncio.wait_for()给protocol.close()添加5秒超时,避免因连接失效导致的无限挂起。 - 优化异常处理:合并同类异常捕获,同时在意外异常时执行重试逻辑(直到达到最大次数)。
- 修复主函数阻塞问题:将
event_loop.run_forever()替换为等待所有pending任务完成,避免无限阻塞。
内容的提问来源于stack exchange,提问作者AVarf
相关产品推荐
相关产品推荐

