Celery任务中如何使用db对象及session属性?是否需调用db.session.remove?
Hey there! Let's walk through your questions about working with Flask-SQLAlchemy's db object and sessions in Celery tasks, especially when dealing with the --maxtaskperchild parameter.
How to use the db object and its session in Celery tasks?
First off, remember that Celery workers run as separate processes from your Flask app. When you initialize db = SQLAlchemy(app), that db instance is tied to your Flask app's context. Celery tasks don't automatically inherit this context, so you need to explicitly activate it before using db or any model queries.
Here's how to adjust your task code to include the context:
from project import celery, db, app from project.models import Model @celery.task() def task_1(): # Activate the Flask app context for this task with app.app_context(): model = Model.query.get(0) model.arg = 1 db.session.commit()
Without this context, you might run into errors like "No application found. Either work inside a view function or push an application context."
Do I need to call db.session.remove() after tasks finish?
Absolutely, and this becomes even more critical when you use the --maxtaskperchild parameter.
SQLAlchemy sessions are thread-local, and Celery workers reuse threads (or processes, when using --maxtaskperchild) for multiple tasks. If you don't clean up the session after a task, leftover connections or state can cause issues like:
- Stale data in subsequent tasks
- Database connection leaks
- Errors when worker processes are restarted (which
--maxtaskperchildtriggers)
To fix this, add a finally block to ensure the session is always removed, even if the task throws an error:
@celery.task() def task_1(): with app.app_context(): try: model = Model.query.get(0) model.arg = 1 db.session.commit() finally: db.session.remove()
A cleaner approach: Use Celery signals
If you have multiple tasks, you can avoid repeating the context setup/teardown code by using Celery's task_prerun and task_postrun signals. These run automatically before and after every task:
from project import celery, db, app from celery.signals import task_prerun, task_postrun @task_prerun.connect def push_app_context(sender=None, **kwargs): app.app_context().push() @task_postrun.connect def cleanup_session(sender=None, **kwargs): db.session.remove() # Now your tasks can skip manual context handling @celery.task() def task_1(): model = Model.query.get(0) model.arg = 1 db.session.commit()
Why does --maxtaskperchild cause issues?
When you set --maxtaskperchild=N, Celery worker processes are destroyed and restarted after executing N tasks. If your tasks don't clean up sessions, the dying process might leave open database connections in an inconsistent state. The next worker process that picks up a task could inherit these stale connections, leading to unexpected errors like "Lost connection to MySQL server" or transaction conflicts.
Calling db.session.remove() ensures that each task's session is properly closed before the worker process is potentially restarted.
内容的提问来源于stack exchange,提问作者Maxime De Bruyn

