Flask后端Python Asyncio后台定时任务启停报错排查与修复
Flask后台定时任务启动/停止触发500错误的修复方案
问题场景
开发Python Flask后端应用,需实现后台定时任务(如每分钟自动更新数据),但调用启动/停止接口时出现以下问题:
- 调用
AsyncPeriodictUtils.start_task():后端返回500错误,但定时任务实际已启动 - 调用
AsyncPeriodictUtils.stop_task():后端返回500错误,且定时任务无法停止
相关代码
view_func.py
from flask import Blueprint, request from async_periodic_job import AsyncPeriodictUtils infras_bp = Blueprint("infras", __name__) @infras_bp.route("/infras/autosync", methods=["PUT"]) def auto_sync(): args = request.args check_required_params(args, ["autosnyc"]) autosnyc = args.get("autosnyc", default="false").lower() == "true" if autosnyc: AsyncPeriodictUtils.start_task() else: AsyncPeriodictUtils.stop_task() return "success", 200
async_periodic_job.py
import asyncio from typing import Callable from utils import logger SECOND = 1 MINUTE = 60 * SECOND logger.debug(f"Import {__file__}") periodic_jobs = dict() task_instance = None # Get default event loop in main thread loop = asyncio.get_event_loop() class AsyncPeriodictUtils: @staticmethod async def run_jobs() -> None: while True: await asyncio.sleep(10) logger.info(f"Called run_jobs periodicly.") logger.info(f"periodic_jobs: {periodic_jobs.keys()}") for func_name, function in periodic_jobs.items(): function() logger.info(f"Called function '{func_name}' periodicly.") @classmethod def create_task(cls) -> None: global task_instance, loop task_instance = loop.create_task(cls.run_jobs()) @staticmethod async def cancel_task() -> None: global task_instance if task_instance: task_instance.cancel() try: await task_instance except asyncio.CancelledError: logger.info("Async periodic task has been cancelled.") task_instance = None else: logger.warning("Async periodic task has not been started yet.") @classmethod def start_task(cls) -> None: cls.create_task() global loop try: loop.run_until_complete(task_instance) except asyncio.CancelledError: pass logger.info("Async Periodic jobs launched.") @classmethod def stop_task(cls) -> None: global loop try: loop.run_until_complete(cls.cancel_task()) except asyncio.CancelledError: pass logger.info("Async Periodic jobs terminated.") @classmethod def add_job(cls, function: Callable) -> None: if function.__name__ in periodic_jobs: return periodic_jobs.update({function.__name__: function}) logger.info(f"Added function {function.__name__} in periodic jobs.") global task_instance if not task_instance: logger.info(f"Auto enable periodic jobs.") cls.start_task() @classmethod def remove_job(cls, function: Callable) -> None: if function.__name__ not in periodic_jobs: logger.warning( f"function {function.__name__} not in periodic jobs.") return periodic_jobs.pop(function.__name__) logger.info(f"Removed function {function.__name__} in periodic jobs.") if not periodic_jobs: logger.info(f"Periodic jobs list clear, auto terminate.") cls.stop_task()
错误日志
启动任务时的500错误日志
[2023-10-13-06:06:19][DAAT][ERROR][__init__.py] [Error] stack: [2023-10-13-06:06:19][DAAT][ERROR][__init__.py] Traceback (most recent call last): File "/opt/conda/lib/python3.10/site-packages/flask/app.py", line 1820, in full_dispatch_request rv = self.dispatch_request() File "/opt/conda/lib/python3.10/site-packages/flask/app.py", line 1796, in dispatch_request return self.ensure_sync(self.view_functions[rule.endpoint])(**view_args) File "/opt/conda/lib/python3.10/site-packages/daat-1.0.0-py3.10.egg/daat/app/routes/infras.py", line 134, in auto_sync AsyncPeriodictUtils.start_task() File "/opt/conda/lib/python3.10/site-packages/daat-1.0.0-py3.10.egg/daat/infras/async_periodic_job.py", line 53, in start_task loop.run_until_complete(task_instance) File "/opt/conda/lib/python3.10/asyncio/base_events.py", line 625, in run_until_complete self._check_running() File "/opt/conda/lib/python3.10/asyncio/base_events.py", line 584, in _check_running raise RuntimeError('This event loop is already running') RuntimeError: This event loop is already running [2023-10-13-06:06:19][werkzeug][INFO][_internal.py] 172.23.0.3 - - [13/Oct/2023 06:06:19] "PUT /infras/autosync?autosnyc=true HTTP/1.1" 500
启动后任务运行日志
[2023-10-13-06:06:23][DAAT][INFO][async_periodic_job.py] Called run_jobs periodicly. [2023-10-13-06:06:23][DAAT][INFO][async_periodic_job.py] periodic_jobs: dict_keys([]) [2023-10-13-06:06:33][DAAT][INFO][async_periodic_job.py] Called run_jobs periodicly. [2023-10-13-06:06:33][DAAT][INFO][async_periodic_job.py] periodic_jobs: dict_keys([])
停止任务时的500错误日志
[2023-10-13-06:06:35][DAAT][ERROR][__init__.py] [Error] stack: [2023-10-13-06:06:35][DAAT][ERROR][__init__.py] Traceback (most recent call last): File "/opt/conda/lib/python3.10/site-packages/flask/app.py", line 1820, in full_dispatch_request rv = self.dispatch_request() File "/opt/conda/lib/python3.10/site-packages/flask/app.py", line 1796, in dispatch_request return self.ensure_sync(self.view_functions[rule.endpoint])(**view_args) File "/opt/conda/lib/python3.10/site-packages/daat-1.0.0-py3.10.egg/daat/app/routes/infras.py", line 136, in auto_sync AsyncPeriodictUtils.stop_task() File "/opt/conda/lib/python3.10/site-packages/daat-1.0.0-py3.10.egg/daat/infras/async_periodic_job.py", line 62, in stop_task loop.run_until_complete(cls.cancel_task()) File "/opt/conda/lib/python3.10/asyncio/base_events.py", line 625, in run_until_complete self._check_running() File "/opt/conda/lib/python3.10/asyncio/base_events.py", line 584, in _check_running raise RuntimeError('This event loop is already running') RuntimeError: This event loop is already running /opt/conda/lib/python3.10/site-packages/flask/app.py:1822: RuntimeWarning: coroutine 'AsyncPeriodictUtils.cancel_task' was never awaited rv = self.handle_user_exception(e) RuntimeWarning: Enable tracemalloc to get the object allocation traceback [2023-10-13-06:06:35][werkzeug][INFO][_internal.py] 172.23.0.3 - - [13/Oct/2023 06:06:35] "PUT /infras/autosync?autosnyc=false HTTP/1.1" 500 -
停止后任务仍运行的日志
[2023-10-13-06:06:43][DAAT][INFO][async_periodic_job.py] Called run_jobs periodicly. [2023-10-13-06:06:43][DAAT][INFO][async_periodic_job.py] periodic_jobs: dict_keys([]) [2023-10-13-06:06:53][DAAT][INFO][async_periodic_job.py] Called run_jobs periodicly. [2023-10-13-06:06:53][DAAT][INFO][async_periodic_job.py] periodic_jobs: dict_keys([])
预期结果
- 调用
AsyncPeriodictUtils.start_task():成功启动定时任务,后端返回200,无500错误 - 调用
AsyncPeriodictUtils.stop_task():成功终止定时任务,后端返回200,无500错误
解决方案
问题根源
- Flask运行时,主线程的事件循环已经处于运行状态,调用
loop.run_until_complete()会触发RuntimeError: This event loop is already running,因为该方法要求事件循环未启动,且会阻塞等待任务完成(但定时任务是无限循环,无法完成)。 - 停止任务时,直接调用
loop.run_until_complete(cls.cancel_task())同样触发上述错误,且协程未被正确执行,导致任务无法停止。
修改后的async_periodic_job.py关键代码
import asyncio from typing import Callable from utils import logger SECOND = 1 MINUTE = 60 * SECOND logger.debug(f"Import {__file__}") periodic_jobs = dict() task_instance = None # Get default event loop in main thread loop = asyncio.get_event_loop() class AsyncPeriodictUtils: @staticmethod async def run_jobs() -> None: while True: await asyncio.sleep(10) logger.info(f"Called run_jobs periodicly.") logger.info(f"periodic_jobs: {periodic_jobs.keys()}") for func_name, function in periodic_jobs.items(): function() logger.info(f"Called function '{func_name}' periodicly.") @classmethod def create_task(cls) -> None: global task_instance, loop # 防止重复创建任务 if task_instance is not None and not task_instance.done(): logger.warning("Async periodic task is already running.") return task_instance = loop.create_task(cls.run_jobs()) @staticmethod async def cancel_task() -> None: global task_instance if task_instance: task_instance.cancel() try: await task_instance except asyncio.CancelledError: logger.info("Async periodic task has been cancelled.") task_instance = None else: logger.warning("Async periodic task has not been started yet.") @classmethod def start_task(cls) -> None: cls.create_task() logger.info("Async Periodic jobs launched.") @classmethod def stop_task(cls) -> None: global loop # 将cancel_task协程提交到已有事件循环执行,而非阻塞等待 loop.create_task(cls.cancel_task()) logger.info("Async Periodic jobs termination requested.") # 其余add_job、remove_job方法保持不变 @classmethod def add_job(cls, function: Callable) -> None: if function.__name__ in periodic_jobs: return periodic_jobs.update({function.__name__: function}) logger.info(f"Added function {function.__name__} in periodic jobs.") global task_instance if not task_instance or (task_instance and task_instance.done()): logger.info(f"Auto enable periodic jobs.") cls.start_task() @classmethod def remove_job(cls, function: Callable) -> None: if function.__name__ not in periodic_jobs: logger.warning( f"function {function.__name__} not in periodic jobs.") return periodic_jobs.pop(function.__name__) logger.info(f"Removed function {function.__name__} in periodic jobs.") if not periodic_jobs: logger.info(f"Periodic jobs list clear, auto terminate.") cls.stop_task()
修改说明
- start_task方法:移除
loop.run_until_complete(task_instance)调用,因为loop.create_task()已经将定时任务加入到运行中的事件循环,无需阻塞等待任务完成。同时在create_task中增加重复启动检查,避免创建多个相同任务。 - stop_task方法:用
loop.create_task(cls.cancel_task())替代loop.run_until_complete(cls.cancel_task()),将取消任务的协程提交到已有事件循环异步执行,不会阻塞请求线程,也不会触发事件循环已运行的错误。 - add_job方法:更新任务存在性判断,确保只有在任务未启动或已完成时才自动启动任务。
验证效果
修改后调用接口:
- 启动任务:后端返回200,定时任务正常运行,无500错误
- 停止任务:后端返回200,定时任务停止运行,无500错误
内容的提问来源于stack exchange,提问作者stevezkw
相关产品推荐
相关产品推荐

