【发布时间】:2018-08-15 16:05:07
【问题描述】:
我正在尝试解决此错误:RuntimeError: Cannot close a running event loop 在我的异步进程中。我相信它正在发生,因为在任务仍然挂起时出现故障,然后我尝试关闭事件循环。我想我需要在关闭事件循环之前等待剩余的响应,但我不确定如何在我的具体情况下正确完成。
def start_job(self):
if self.auth_expire_timestamp < get_timestamp():
api_obj = api_handler.Api('Api Name', self.dbObj)
self.api_auth_resp = api_obj.get_auth_response()
self.api_attr = api_obj.get_attributes()
try:
self.queue_manager(self.do_stuff(json_data))
except aiohttp.ServerDisconnectedError as e:
logging.info("Reconnecting...")
api_obj = api_handler.Api('API Name', self.dbObj)
self.api_auth_resp = api_obj.get_auth_response()
self.api_attr = api_obj.get_attributes()
self.run_eligibility()
async def do_stuff(self, data):
tasks = []
async with aiohttp.ClientSession() as session:
for row in data:
task = asyncio.ensure_future(self.async_post('url', session, row))
tasks.append(task)
result = await asyncio.gather(*tasks)
self.load_results(result)
def queue_manager(self, method):
self.loop = asyncio.get_event_loop()
future = asyncio.ensure_future(method)
self.loop.run_until_complete(future)
async def async_post(self, resource, session, data):
async with session.post(self.api_attr.api_endpoint + resource, headers=self.headers, data=data) as response:
resp = []
try:
headers = response.headers['foo']
content = await response.read()
resp.append(headers)
resp.append(content)
except KeyError as e:
logging.error('KeyError at async_post response')
logging.error(e)
return resp
def shutdown(self):
//need to do something here to await the remaining tasks and then I need to re-start a new event loop, which i think i can do, just don't know how to appropriately stop the current one.
self.loop.close()
return True
我如何处理错误并正确关闭事件循环,以便我可以启动一个新的并基本上重新启动整个程序并继续。
编辑:
这就是我现在正在尝试的,基于this SO answer。不幸的是,这个错误很少发生,所以除非我能强制它,否则我将不得不等待,看看它是否有效。在我的queue_manager 方法中,我将其更改为:
try:
self.loop.run_until_complete(future)
except Exception as e:
future.cancel()
self.loop.run_until_complete(future)
future.exception()
更新:
我摆脱了shutdown() 方法并将其添加到我的queue_manager() 方法中,它似乎可以正常工作:
try:
self.loop.run_until_complete(future)
except Exception as e:
future.cancel()
self.check_in_records()
self.reconnect()
self.start_job()
future.exception()
【问题讨论】:
-
shutdown是从哪里调用的,为什么要尝试close事件循环?一个 asyncio 程序通常由一个事件循环实例在其整个生命周期内提供服务。 -
问题是我调用的 API 将与待处理的任务断开连接,我试图“重新启动”而不会使整个应用程序崩溃。我为我最近添加的内容添加了更新。它似乎有效,但我愿意接受反馈。
-
最后需要
future.exception()吗?run_until_complete似乎正确地发现了异常。另外,取消显然已经完成的未来有什么意义(run_until_completeraise 见证了这一点)? -
这个答案 - stackoverflow.com/a/30766124/4113027 似乎表明我需要
future.exception()位。至于你的另一个问题,我正在取消未来,因为在这种情况下还有剩余的任务待处理,我想基本上删除这些任务并开始一个新的事件循环。在这种情况下,服务器在所有任务完成之前就断开了连接,因此那些剩余的任务无论如何都不会带回任何数据......也许我没有正确处理......? -
在您的代码中,
cancel()和exception()都是不必要的。cancel()因为未来已经完成(有一个例外),run_until_complete退出的事实证明了这一点。请参阅this code 以获取功能等效的示例,该示例在没有警告的情况下运行。
标签: python python-3.6 python-asyncio aiohttp