如何将Flask模型中基于RQ的launch_task方法改为使用Celery?
用Celery替代RQ实现Flask用户模型的后台任务调度
在Miguel Grinberg的《Flask Mega-Tutorial》第22章中,作者为User对象模型实现了通用的launch_task方法,基于RQ处理后台任务。以下是将该方法转换为Celery实现的完整步骤:
原RQ实现参考代码
User类中的launch_task方法
class User(UserMixin, db.Model): # ... 其他字段和方法 ... def launch_task(self, name, description, *args, **kwargs): rq_job = current_app.task_queue.enqueue(f'app.tasks.{name}', self.id, *args, **kwargs) task = Task(id=rq_job.get_id(), name=name, description=description, user=self) db.session.add(task) return task
任务调用示例
@bp.route('/export_posts') @login_required def export_posts(): if current_user.get_task_in_progress('export_posts'): flash(_('An export task is currently in progress')) else: current_user.launch_task('export_posts', _('Exporting posts...')) db.session.commit() return redirect(url_for('main.user', username=current_user.username))
tasks.py中的任务函数
def export_posts(user_id): try: # 从数据库读取用户文章 # 给用户发送包含数据的邮件 except Exception: # 处理意外错误 finally: # 清理操作
转换为Celery实现的步骤
1. 完成Celery基础配置
首先在Flask应用中初始化Celery实例,配置好消息代理和结果后端:
# app/__init__.py from flask import Flask from celery import Celery def create_app(): app = Flask(__name__) # Celery配置示例(使用Redis作为代理和结果后端) app.config['CELERY_BROKER_URL'] = 'redis://localhost:6379/0' app.config['CELERY_RESULT_BACKEND'] = 'redis://localhost:6379/0' celery = Celery(app.name, broker=app.config['CELERY_BROKER_URL']) celery.conf.update(app.config) app.celery = celery # 将Celery实例挂载到Flask应用上 # ... 其他应用初始化逻辑 ... return app
2. 修改User类的launch_task方法
替换RQ的任务入队逻辑为Celery的任务发送方式,获取Celery任务ID并关联到Task模型:
class User(UserMixin, db.Model): # ... 其他字段和方法 ... def launch_task(self, name, description, *args, **kwargs): # 用Celery发送任务,指定任务路径、用户ID及参数 celery_task = current_app.celery.send_task( f'app.tasks.{name}', args=(self.id, *args), kwargs=kwargs ) # 创建Task记录,使用Celery任务ID作为主键 task = Task( id=celery_task.id, name=name, description=description, user=self ) db.session.add(task) return task
3. 改造tasks.py中的任务函数
给任务添加Celery装饰器,支持任务状态更新:
from app import celery from app.models import Task, User, db @celery.task(bind=True) def export_posts(self, user_id): try: # 业务逻辑:读取用户文章、发送邮件等 user = User.query.get(user_id) # ... 具体实现代码 ... except Exception as e: # 更新任务状态为失败,记录错误信息 self.update_state(state='FAILURE', meta={'error': str(e)}) # 同步到Task模型 task = Task.query.get(self.request.id) if task: task.complete = True task.error = str(e) db.session.commit() raise finally: # 任务完成后更新Task模型状态 task = Task.query.get(self.request.id) if task: task.complete = True db.session.commit() # 其他清理操作
使用bind=True可以让任务函数获取到任务实例,方便更新状态和关联数据库记录。
4. 任务状态检查逻辑复用
原代码中的get_task_in_progress方法无需修改——只要Task模型包含complete字段,就能继续通过查询未完成的同名任务来判断是否有任务在执行。
5. 可选:用Celery信号自动更新任务状态
如果不想在每个任务里手动更新Task模型,可以用Celery的信号统一处理:
# app/signals.py from celery.signals import task_success, task_failure from app.models import Task, db @task_success.connect def update_task_success(sender=None, result=None, **kwargs): task = Task.query.get(sender.request.id) if task: task.complete = True db.session.commit() @task_failure.connect def update_task_failure(sender=None, exception=None, **kwargs): task = Task.query.get(sender.request.id) if task: task.complete = True task.error = str(exception) db.session.commit()
记得在应用初始化时导入这个信号文件,确保信号生效。
内容的提问来源于stack exchange,提问作者setty
相关产品推荐
相关产品推荐

