最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Celery定時(shí)任務(wù)組件之Django+Celery項(xiàng)目實(shí)戰(zhàn)教程

 更新時(shí)間:2025年07月02日 11:08:24   作者:alden_ygq  
這篇文章主要介紹了Celery定時(shí)任務(wù)組件之Django+Celery項(xiàng)目實(shí)戰(zhàn),具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教

一、項(xiàng)目初始化

1. 創(chuàng)建虛擬環(huán)境并安裝依賴

# 創(chuàng)建虛擬環(huán)境
python3 -m venv myenv
source myenv/bin/activate

# 安裝依賴
pip install django celery redis django-celery-beat

2. 創(chuàng)建 Django 項(xiàng)目和應(yīng)用

# 創(chuàng)建項(xiàng)目
django-admin startproject task_manager
cd task_manager

# 創(chuàng)建應(yīng)用
python manage.py startapp tasks

3. 配置項(xiàng)目(task_manager/settings.py)

INSTALLED_APPS = [
    # ...
    'django_celery_beat',
    'django_celery_results',
    'tasks',
]

# 數(shù)據(jù)庫配置
DATABASES = {
    'default': {
        'ENGINE': 'django.db.backends.sqlite3',
        'NAME': BASE_DIR / 'db.sqlite3',
    }
}

# Celery配置
CELERY_BROKER_URL = 'redis://localhost:6379/0'
CELERY_RESULT_BACKEND = 'django-db'  # 使用django-celery-results存儲(chǔ)結(jié)果
CELERY_ACCEPT_CONTENT = ['json']
CELERY_TASK_SERIALIZER = 'json'
CELERY_RESULT_SERIALIZER = 'json'
CELERY_TIMEZONE = 'Asia/Shanghai'

二、Celery 集成配置

1. 創(chuàng)建 Celery 應(yīng)用(task_manager/celery.py)

from __future__ import absolute_import, unicode_literals
import os
from celery import Celery

os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'task_manager.settings')

app = Celery('task_manager')
app.config_from_object('django.conf:settings', namespace='CELERY')
app.autodiscover_tasks()

@app.task(bind=True)
def debug_task(self):
    print(f'Request: {self.request!r}')

2. 初始化 Celery(task_manager/__init__.py)

from __future__ import absolute_import, unicode_literals
from .celery import app as celery_app

__all__ = ('celery_app',)

三、Model 開發(fā)

創(chuàng)建任務(wù)模型(tasks/models.py)

from django.db import models
from django.utils import timezone

class ScheduledTask(models.Model):
    TASK_TYPES = (
        ('periodic', '周期性任務(wù)'),
        ('one_time', '一次性任務(wù)'),
    )
    
    name = models.CharField('任務(wù)名稱', max_length=100)
    task_type = models.CharField('任務(wù)類型', max_length=20, choices=TASK_TYPES)
    task_function = models.CharField('任務(wù)函數(shù)', max_length=200)
    cron_expression = models.CharField('Cron表達(dá)式', max_length=100, blank=True, null=True)
    interval_seconds = models.IntegerField('間隔秒數(shù)', blank=True, null=True)
    next_run_time = models.DateTimeField('下次執(zhí)行時(shí)間', blank=True, null=True)
    is_active = models.BooleanField('是否激活', default=True)
    created_at = models.DateTimeField('創(chuàng)建時(shí)間', auto_now_add=True)
    updated_at = models.DateTimeField('更新時(shí)間', auto_now=True)
    
    def __str__(self):
        return self.name
    
    class Meta:
        verbose_name = '定時(shí)任務(wù)'
        verbose_name_plural = '定時(shí)任務(wù)列表'

class TaskExecutionLog(models.Model):
    task = models.ForeignKey(ScheduledTask, on_delete=models.CASCADE, related_name='logs')
    execution_time = models.DateTimeField('執(zhí)行時(shí)間', auto_now_add=True)
    status = models.CharField('執(zhí)行狀態(tài)', max_length=20, choices=(
        ('success', '成功'),
        ('failed', '失敗'),
    ))
    result = models.TextField('執(zhí)行結(jié)果', blank=True, null=True)
    error_message = models.TextField('錯(cuò)誤信息', blank=True, null=True)
    
    def __str__(self):
        return f"{self.task.name} - {self.execution_time}"
    
    class Meta:
        verbose_name = '任務(wù)執(zhí)行日志'
        verbose_name_plural = '任務(wù)執(zhí)行日志列表'

遷移數(shù)據(jù)庫

python manage.py makemigrations
python manage.py migrate

四、接口開發(fā)

1. 創(chuàng)建序列化器(tasks/serializers.py)

from rest_framework import serializers
from .models import ScheduledTask, TaskExecutionLog

class ScheduledTaskSerializer(serializers.ModelSerializer):
    class Meta:
        model = ScheduledTask
        fields = '__all__'

class TaskExecutionLogSerializer(serializers.ModelSerializer):
    class Meta:
        model = TaskExecutionLog
        fields = '__all__'

2. 創(chuàng)建視圖集(tasks/views.py)

from rest_framework import viewsets, status
from rest_framework.response import Response
from .models import ScheduledTask, TaskExecutionLog
from .serializers import ScheduledTaskSerializer, TaskExecutionLogSerializer
from celery import current_app
from django_celery_beat.models import PeriodicTask, IntervalSchedule, CrontabSchedule
import json

class ScheduledTaskViewSet(viewsets.ModelViewSet):
    queryset = ScheduledTask.objects.all()
    serializer_class = ScheduledTaskSerializer
    
    def create(self, request, *args, **kwargs):
        serializer = self.get_serializer(data=request.data)
        serializer.is_valid(raise_exception=True)
        
        # 創(chuàng)建Celery定時(shí)任務(wù)
        task = serializer.save()
        self._create_celery_task(task)
        
        headers = self.get_success_headers(serializer.data)
        return Response(serializer.data, status=status.HTTP_201_CREATED, headers=headers)
    
    def update(self, request, *args, **kwargs):
        partial = kwargs.pop('partial', False)
        instance = self.get_object()
        serializer = self.get_serializer(instance, data=request.data, partial=partial)
        serializer.is_valid(raise_exception=True)
        
        # 更新Celery定時(shí)任務(wù)
        task = serializer.save()
        self._update_celery_task(task)
        
        return Response(serializer.data)
    
    def destroy(self, request, *args, **kwargs):
        instance = self.get_object()
        
        # 刪除Celery定時(shí)任務(wù)
        self._delete_celery_task(instance)
        
        self.perform_destroy(instance)
        return Response(status=status.HTTP_204_NO_CONTENT)
    
    def _create_celery_task(self, task):
        if task.task_type == 'periodic':
            # 創(chuàng)建間隔調(diào)度
            schedule, _ = IntervalSchedule.objects.get_or_create(
                every=task.interval_seconds,
                period=IntervalSchedule.SECONDS,
            )
            PeriodicTask.objects.create(
                interval=schedule,
                name=task.name,
                task=task.task_function,
                enabled=task.is_active,
                args=json.dumps([]),
                kwargs=json.dumps({}),
            )
        elif task.task_type == 'one_time':
            # 一次性任務(wù)使用ETA
            pass
    
    def _update_celery_task(self, task):
        try:
            periodic_task = PeriodicTask.objects.get(name=task.name)
            if task.task_type == 'periodic':
                schedule, _ = IntervalSchedule.objects.get_or_create(
                    every=task.interval_seconds,
                    period=IntervalSchedule.SECONDS,
                )
                periodic_task.interval = schedule
            periodic_task.enabled = task.is_active
            periodic_task.save()
        except PeriodicTask.DoesNotExist:
            self._create_celery_task(task)
    
    def _delete_celery_task(self, task):
        try:
            periodic_task = PeriodicTask.objects.get(name=task.name)
            periodic_task.delete()
        except PeriodicTask.DoesNotExist:
            pass

class TaskExecutionLogViewSet(viewsets.ReadOnlyModelViewSet):
    queryset = TaskExecutionLog.objects.all()
    serializer_class = TaskExecutionLogSerializer

3. 配置 URL(tasks/urls.py)

from django.urls import include, path
from rest_framework import routers
from .views import ScheduledTaskViewSet, TaskExecutionLogViewSet

router = routers.DefaultRouter()
router.register(r'tasks', ScheduledTaskViewSet)
router.register(r'logs', TaskExecutionLogViewSet)

urlpatterns = [
    path('', include(router.urls)),
]

4. 項(xiàng)目 URL 配置(task_manager/urls.py)

from django.contrib import admin
from django.urls import path, include

urlpatterns = [
    path('admin/', admin.site.urls),
    path('api/', include('tasks.urls')),
]

五、創(chuàng)建示例任務(wù)

定義任務(wù)函數(shù)(tasks/tasks.py)

from celery import shared_task
from .models import ScheduledTask, TaskExecutionLog
import logging

logger = logging.getLogger(__name__)

@shared_task(bind=True, autoretry_for=(Exception,), retry_backoff=3, retry_kwargs={'max_retries': 3})
def sample_task(self, task_id):
    try:
        task = ScheduledTask.objects.get(id=task_id)
        
        # 模擬任務(wù)執(zhí)行
        result = f"任務(wù) {task.name} 執(zhí)行成功,時(shí)間:{str(self.request.time_start)}"
        
        # 記錄執(zhí)行日志
        TaskExecutionLog.objects.create(
            task=task,
            status='success',
            result=result
        )
        
        logger.info(f"任務(wù)執(zhí)行成功: {task.name}")
        return result
        
    except Exception as e:
        # 記錄錯(cuò)誤日志
        task = ScheduledTask.objects.get(id=task_id) if ScheduledTask.objects.filter(id=task_id).exists() else None
        if task:
            TaskExecutionLog.objects.create(
                task=task,
                status='failed',
                error_message=str(e)
            )
        logger.error(f"任務(wù)執(zhí)行失敗: {str(e)}")
        raise

六、啟動(dòng)服務(wù)

1. 啟動(dòng) Redis

redis-server

2. 啟動(dòng) Celery Worker

celery -A task_manager worker --loglevel=info --pool=prefork --concurrency=4

3. 啟動(dòng) Celery Beat

celery -A task_manager beat --loglevel=info --scheduler django_celery_beat.schedulers:DatabaseScheduler

4. 啟動(dòng) Django 開發(fā)服務(wù)器

python manage.py runserver

七、API 測(cè)試

1. 創(chuàng)建周期性任務(wù)

curl -X POST http://localhost:8000/api/tasks/ -d '{
    "name": "示例周期性任務(wù)",
    "task_type": "periodic",
    "task_function": "tasks.tasks.sample_task",
    "interval_seconds": 60,
    "is_active": true
}' -H "Content-Type: application/json"

2. 查看任務(wù)列表

curl http://localhost:8000/api/tasks/

3. 查看執(zhí)行日志

curl http://localhost:8000/api/logs/

項(xiàng)目結(jié)構(gòu)

task_manager/
├── task_manager/
│   ├── __init__.py
│   ├── celery.py
│   ├── settings.py
│   ├── urls.py
│   └── wsgi.py
├── tasks/
│   ├── migrations/
│   ├── __init__.py
│   ├── admin.py
│   ├── apps.py
│   ├── models.py
│   ├── serializers.py
│   ├── tasks.py
│   ├── urls.py
│   └── views.py
├── manage.py
└── db.sqlite3

關(guān)鍵特性說明

  • 動(dòng)態(tài)任務(wù)管理:通過 API 創(chuàng)建 / 更新 / 刪除定時(shí)任務(wù)
  • 任務(wù)執(zhí)行記錄:自動(dòng)記錄任務(wù)執(zhí)行結(jié)果和狀態(tài)
  • 失敗重試機(jī)制:任務(wù)失敗時(shí)自動(dòng)重試(最多 3 次)
  • 多種調(diào)度方式:支持周期性任務(wù)和一次性任務(wù)
  • 可視化管理:通過 Django Admin 界面管理定時(shí)任務(wù)

擴(kuò)展建議

  • 添加任務(wù)參數(shù)支持,允許在創(chuàng)建任務(wù)時(shí)傳遞參數(shù)
  • 實(shí)現(xiàn)任務(wù)暫停 / 恢復(fù)功能
  • 添加任務(wù)優(yōu)先級(jí)隊(duì)列配置
  • 集成監(jiān)控系統(tǒng)(如 Prometheus+Grafana)
  • 實(shí)現(xiàn)任務(wù)執(zhí)行結(jié)果的異步通知(郵件、短信等)

這個(gè)實(shí)現(xiàn)提供了一個(gè)完整的 Django+Celery 定時(shí)任務(wù)系統(tǒng),支持動(dòng)態(tài)管理和監(jiān)控,可直接用于生產(chǎn)環(huán)境。

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

最新評(píng)論

门源| 敦煌市| 左权县| 华安县| 渑池县| 晴隆县| 漠河县| 六枝特区| 石家庄市| 广昌县| 淮北市| 灵宝市| 交城县| 诸城市| 永吉县| 望奎县| 广州市| 水富县| 马鞍山市| 教育| 上饶市| 巩留县| 循化| 上犹县| 延川县| 灵璧县| 若尔盖县| 固原市| 彰武县| 方正县| 昌邑市| 新郑市| 万山特区| 镇雄县| 布尔津县| 博野县| 武威市| 巴马| 景宁| 昌平区| 绥江县|