You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 15:53:08