Django中扩展Celery存储所有任务状态至PostgreSQL方案咨询
扩展Celery-Results存储所有任务状态(含RECEIVED)的实现方案
1. 自定义数据库结果后端
默认的django-db结果后端只存储完成状态的任务,直接继承它重写状态存储逻辑,让所有状态变更都写入数据库:
# your_project/celery_backends.py from celery.backends.database import DatabaseBackend from celery.backends.database.models import TaskResult from django.db import transaction from django.utils import timezone class FullStateDatabaseBackend(DatabaseBackend): def _store_result(self, task_id, result, state, traceback=None, request=None, **kwargs): # 无论任务处于什么状态,都创建或更新数据库记录 now = timezone.now() with transaction.atomic(): obj, created = TaskResult.objects.get_or_create( task_id=task_id, defaults={ 'status': state, 'result': self.encode(result), 'traceback': traceback, 'meta': self.encode(kwargs.get('meta', {})), } ) if not created: obj.status = state obj.result = self.encode(result) obj.traceback = traceback obj.meta = self.encode(kwargs.get('meta', {})) # 可选:根据状态更新时间字段 if state == 'RECEIVED': obj.date_created = now elif state == 'STARTED': obj.date_done = now obj.save(update_fields=['status', 'result', 'traceback', 'meta', 'date_created', 'date_done']) return obj
2. 配置Celery使用自定义后端
在项目的celery.py里修改配置,指定自定义结果后端:
# your_project/celery.py from celery import Celery import os os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'your_project.settings') app = Celery('your_project') app.config_from_object('django.conf:settings', namespace='CELERY') # 替换为自定义后端的路径 app.conf.result_backend = 'your_project.celery_backends.FullStateDatabaseBackend'
3. 用Celery信号捕获RECEIVED状态
部分中间状态(比如RECEIVED)不会触发默认的结果存储逻辑,需要靠Celery信号主动捕获记录。在Django应用的apps.py里注册信号处理器:
# your_app/apps.py from django.apps import AppConfig from celery import signals from celery.backends.database.models import TaskResult from django.db import transaction from django.utils import timezone class YourAppConfig(AppConfig): default_auto_field = 'django.db.models.BigAutoField' name = 'your_app' def ready(self): # 捕获任务被接收的信号 @signals.task_received.connect def log_task_received(sender=None, request=None, **kwargs): if not request or not request.id: return with transaction.atomic(): TaskResult.objects.update_or_create( task_id=request.id, defaults={ 'status': 'RECEIVED', 'task_name': request.task, 'date_created': timezone.now(), 'meta': request.encode({'args': request.args, 'kwargs': request.kwargs}) } ) # 捕获任务启动的信号(可选,确保STARTED状态被记录) @signals.task_started.connect def log_task_started(sender=None, task_id=None, **kwargs): if not task_id: return with transaction.atomic(): obj, _ = TaskResult.objects.get_or_create(task_id=task_id) obj.status = 'STARTED' obj.date_done = timezone.now() obj.save(update_fields=['status', 'date_done']) # 按需添加task_revoked、task_failed等其他状态的信号处理
4. 扩展TaskResult模型(可选)
如果需要存储更多元数据(比如单独的任务接收、启动时间字段),可以自定义模型继承原TaskResult:
# your_app/models.py from celery.backends.database.models import TaskResult as BaseTaskResult from django.db import models class ExtendedTaskResult(BaseTaskResult): received_at = models.DateTimeField(null=True, blank=True) started_at = models.DateTimeField(null=True, blank=True) class Meta: db_table = 'extended_celery_taskmeta'
回到自定义后端,替换使用扩展后的模型:
# 在FullStateDatabaseBackend中添加指定模型的代码 class FullStateDatabaseBackend(DatabaseBackend): _model = ExtendedTaskResult # 指定自定义模型 def _store_result(self, task_id, result, state, traceback=None, request=None, **kwargs): now = timezone.now() with transaction.atomic(): obj, created = self._model.objects.get_or_create( task_id=task_id, defaults={ 'status': state, 'result': self.encode(result), 'traceback': traceback, 'meta': self.encode(kwargs.get('meta', {})), 'received_at': now if state == 'RECEIVED' else None, 'started_at': now if state == 'STARTED' else None, } ) if not created: obj.status = state obj.result = self.encode(result) obj.traceback = traceback obj.meta = self.encode(kwargs.get('meta', {})) if state == 'RECEIVED': obj.received_at = now elif state == 'STARTED': obj.started_at = now obj.save(update_fields=['status', 'result', 'traceback', 'meta', 'received_at', 'started_at']) return obj
5. 迁移数据库并验证
- 如果扩展了模型,执行数据库迁移:
python manage.py makemigrations python manage.py migrate
- 重启Celery worker和beat:
celery -A your_project worker -l info celery -A your_project beat -l info
- 触发测试任务,查看数据库对应表(
extended_celery_taskmeta或原celery_taskmeta),确认RECEIVED、STARTED、SUCCESS/FAILURE等状态均已记录。
之后React前端即可通过Django API接口(比如用DRF实现)查询这些数据,完成全状态任务监控。
内容的提问来源于stack exchange,提问作者Adrian
相关产品推荐
相关产品推荐

