如何在Tornado中退出gen.sleep状态并重启刷新工作流?
解决方案:中断Tornado的gen.sleep并重启工作流
核心思路是用可中断的等待逻辑替换原有的gen.sleep,结合事件(Event)实现sleep过程中随时响应中断信号,进而重启工作流。
修正后的完整代码示例
1. 完善RefreshPassword类
import asyncio from tornado import gen, locks, ioloop, web import logging logger = logging.getLogger(__name__) retry_secs = 1800 # 30分钟 SLEEP_SECS = 60 # 示例休眠时间 class RefreshPassword: def __init__(self): self._state = "STOPPED" self._event = locks.Event() def state(self): return self._state def start(self): if self.state() != "ACTIVE": self._state = "ACTIVE" self._event.clear() ioloop.IOLoop.current().add_callback(self._loop) logger.info("Refresh workflow started") return def stop(self): if self.state() == "ACTIVE": self._state = "STOPPED" self._event.set() logger.info("Refresh workflow stopped") async def _interruptible_sleep(self, seconds): # 同时等待休眠完成或事件触发,优先响应先完成的操作 done, _ = await asyncio.wait( [gen.sleep(seconds), self._event.wait()], return_when=asyncio.FIRST_COMPLETED ) # 若由事件触发中断,重置事件状态以支持下次使用 if self._event.is_set(): self._event.clear() async def _loop(self): """Refresh logic.""" while self.state() == "ACTIVE": if self.auth_enabled(): try: password = await _get_metadata() if password: _store_password(password) except Exception as e: logger.exception("Failed to refresh password") await self._interruptible_sleep(retry_secs) else: logger.info("Auth disabled, sleeping...") await self._interruptible_sleep(SLEEP_SECS) logger.info("Refresh loop exited") def auth_enabled(self): # 示例方法,根据实际逻辑实现权限判断 return True async def _get_metadata(): # 示例方法,模拟获取密码逻辑 return "new_password" def _store_password(password): # 示例方法,模拟存储密码逻辑 logger.info(f"Stored new password: {password}")
2. 修正PasswordHandler请求处理器
class PasswordHandler(web.RequestHandler): """Handles HTTP requests""" SUPPORTED_METHODS = ["GET", "POST"] def initialize(self, calendar: RefreshPassword): """Initializes RefreshPassword calendar object. Args: calendar(RefreshPassword): Keeps Password updated. """ self._calendar = calendar async def post(self): msg = "" try: restart = self.get_argument("restart", default="false") if restart.lower() != "true": msg = "Restart requires 'restart=true' parameter" logger.warning(msg) self.set_status(409) self.write(msg) self.finish() return if self._calendar.state() == "ACTIVE": logger.info("Stopping refresh process for restart") self._calendar.stop() # 等待循环退出(可选,确保旧进程完全终止) await gen.sleep(0.1) self._calendar.start() self.set_status(201) msg = "Refresh workflow restarted successfully" self.write(msg) except web.MissingArgumentError as exc: logger.exception("Missing request argument") self.set_status(400) self.write("Missing required argument") finally: self.finish()
关键改进点说明
替换不可中断的
gen.sleep
用asyncio.wait同时监听休眠任务和事件信号,当调用stop()触发事件时,等待逻辑会立即返回,无需等到休眠时间结束。事件状态重置
在_interruptible_sleep中清除事件状态,确保后续等待操作能正常响应新的中断信号。流程逻辑修正
- 修复
start()方法的语法错误(移除多余冒号) - 补充
state()方法实现状态查询 - 完善请求处理器的参数校验和响应逻辑,避免未定义变量错误
- 修复
重启流程保障
在post()方法中,停止旧工作流后短暂等待,确保旧循环完全退出再启动新流程,避免资源竞争。
内容的提问来源于stack exchange,提问作者Nuria
相关产品推荐
相关产品推荐

