Celery定時任務(wù)組件之Django+Celery項(xiàng)目實(shí)戰(zhàn)教程
一、項(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存儲結(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í)行時間', blank=True, null=True)
is_active = models.BooleanField('是否激活', default=True)
created_at = models.DateTimeField('創(chuàng)建時間', auto_now_add=True)
updated_at = models.DateTimeField('更新時間', auto_now=True)
def __str__(self):
return self.name
class Meta:
verbose_name = '定時任務(wù)'
verbose_name_plural = '定時任務(wù)列表'
class TaskExecutionLog(models.Model):
task = models.ForeignKey(ScheduledTask, on_delete=models.CASCADE, related_name='logs')
execution_time = models.DateTimeField('執(zhí)行時間', 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('錯誤信息', 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定時任務(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定時任務(wù)
task = serializer.save()
self._update_celery_task(task)
return Response(serializer.data)
def destroy(self, request, *args, **kwargs):
instance = self.get_object()
# 刪除Celery定時任務(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í)行成功,時間:{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:
# 記錄錯誤日志
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
六、啟動服務(wù)
1. 啟動 Redis
redis-server
2. 啟動 Celery Worker
celery -A task_manager worker --loglevel=info --pool=prefork --concurrency=4
3. 啟動 Celery Beat
celery -A task_manager beat --loglevel=info --scheduler django_celery_beat.schedulers:DatabaseScheduler
4. 啟動 Django 開發(fā)服務(wù)器
python manage.py runserver
七、API 測試
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)鍵特性說明
- 動態(tài)任務(wù)管理:通過 API 創(chuàng)建 / 更新 / 刪除定時任務(wù)
- 任務(wù)執(zhí)行記錄:自動記錄任務(wù)執(zhí)行結(jié)果和狀態(tài)
- 失敗重試機(jī)制:任務(wù)失敗時自動重試(最多 3 次)
- 多種調(diào)度方式:支持周期性任務(wù)和一次性任務(wù)
- 可視化管理:通過 Django Admin 界面管理定時任務(wù)
擴(kuò)展建議
- 添加任務(wù)參數(shù)支持,允許在創(chuàng)建任務(wù)時傳遞參數(shù)
- 實(shí)現(xiàn)任務(wù)暫停 / 恢復(fù)功能
- 添加任務(wù)優(yōu)先級隊(duì)列配置
- 集成監(jiān)控系統(tǒng)(如 Prometheus+Grafana)
- 實(shí)現(xiàn)任務(wù)執(zhí)行結(jié)果的異步通知(郵件、短信等)
這個實(shí)現(xiàn)提供了一個完整的 Django+Celery 定時任務(wù)系統(tǒng),支持動態(tài)管理和監(jiān)控,可直接用于生產(chǎn)環(huán)境。
總結(jié)
以上為個人經(jīng)驗(yàn),希望能給大家一個參考,也希望大家多多支持腳本之家。
相關(guān)文章
計(jì)算機(jī)二級python學(xué)習(xí)教程(1) 教大家如何學(xué)習(xí)python
這篇文章主要為大家詳細(xì)介紹了計(jì)算機(jī)二級python學(xué)習(xí)教程,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2019-05-05
Python面向?qū)ο蠡A(chǔ)入門之設(shè)置對象屬性
這篇文章主要給大家介紹了關(guān)于Python面向?qū)ο蠡A(chǔ)入門之設(shè)置對象屬性的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2018-12-12
Python Social Auth構(gòu)建靈活而強(qiáng)大的社交登錄系統(tǒng)實(shí)例探究
這篇文章主要為大家介紹了Python Social Auth構(gòu)建靈活而強(qiáng)大的社交登錄系統(tǒng)實(shí)例探究,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪2024-01-01
Python正則表達(dá)式語法及re模塊中的常用函數(shù)詳解
這篇文章主要給大家介紹了關(guān)于Python正則表達(dá)式語法及re模塊中常用函數(shù)的相關(guān)資料,正則表達(dá)式是一種強(qiáng)大的字符串處理工具,可以用于匹配、切分、查找和替換等操作,文章詳細(xì)講解了如何使用正則表達(dá)式進(jìn)行匹配,,需要的朋友可以參考下2025-04-04
python進(jìn)程池實(shí)現(xiàn)的多進(jìn)程文件夾copy器完整示例
這篇文章主要介紹了python進(jìn)程池實(shí)現(xiàn)的多進(jìn)程文件夾copy器,結(jié)合完整實(shí)例形式分析了Python基于多進(jìn)程與進(jìn)程池的文件操作相關(guān)實(shí)現(xiàn)技巧,需要的朋友可以參考下2019-11-11
基于Python實(shí)現(xiàn)Windows桌面定時提醒休息程序
這篇文章為大家詳細(xì)主要介紹了如何基于Python實(shí)現(xiàn)簡單的Windows桌面定時提醒休息程序,文中的示例代碼講解詳細(xì),有需要的可以參考一下2024-11-11
Python+Pygame實(shí)戰(zhàn)之英文版猜字游戲的實(shí)現(xiàn)
這篇文章主要為大家介紹了如何利用Python中的Pygame模塊實(shí)現(xiàn)英文版猜單詞游戲,文中的示例代碼講解詳細(xì),對我們學(xué)習(xí)Python游戲開發(fā)有一定幫助,需要的可以參考一下2022-08-08

