Celery 与 Django 集成实战:从零搭建异步任务队列(celery.py、shared_task 与事务安全触发)
发布时间:2026/9/19 7:38:34来源:尧图网络
Celery 与 Django 集成实战从零搭建异步任务队列celery.py、shared_task 与事务安全触发【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery导读本文基于 Celery 官方文档的 Django 集成指南docs/django/first-steps-with-django.rst结合仓库中的完整 Django 示例项目examples/django与底层源码celery/contrib/django/task.py、celery/fixups/django.py系统讲解如何在 Django 项目中引入 Celery 异步任务队列。读完本文你将掌握标准 Django 项目布局下的 Celery 实例定义、shared_task的可复用任务编写方式、在数据库事务提交后再触发任务的正确姿势以及django-celery-results结果后端与 worker 进程的启动与配置。自 Celery 3.1 起Django 支持开箱即用不再需要单独的第三方集成库。本文介绍的是最基本、最标准的 Celery Django 集成方式其 API 与不使用 Django 的用户完全一致因此建议先阅读 First steps with Celery 入门教程再回到本指南。Django 版本兼容性在开始之前先确认 Celery 与 Django 的版本搭配当前仓库主分支Celery 5.5.x支持Django 2.2 LTS 或更新版本若你的 Django 版本低于 2.2请使用Celery 5.2.x若你的 Django 版本低于 1.11请使用Celery 4.4.x。这一约束在源码中同样有体现celery/fixups/django.py 的_verify_django_version会校验django.VERSION (1, 11)并抛出ImproperlyConfigured异常保证集成代码所依赖的 Django API 始终可用。第一步定义 Celery 实例proj/proj/celery.py假设你有一个现代 Django 项目布局proj/ - manage.py - proj/ - __init__.py - settings.py - urls.py推荐的集成方式是新建一个proj/proj/celery.py模块来定义 Celery 实例。仓库中的完整示例见 examples/django/proj/celery.py其核心内容如下import os from celery import Celery # 为 celery 命令行程序设置默认的 Django settings 模块 os.environ.setdefault(DJANGO_SETTINGS_MODULE, proj.settings) app Celery(proj) # 使用字符串意味着 worker 不需要把配置对象序列化给子进程 # - namespaceCELERY 表示所有 celery 相关配置键都应带 CELERY_ 前缀 app.config_from_object(django.conf:settings, namespaceCELERY) # 从所有已注册的 Django app 中加载任务模块 app.autodiscover_tasks() app.task(bindTrue, ignore_resultTrue) def debug_task(self): print(fRequest: {self.request!r})然后在proj/proj/__init__.py中导入这个 app见 examples/django/proj/init.py# 这可以确保 Django 启动时该 app 一定被导入 # 从而使 shared_task 使用这个 app from .celery import app as celery_app __all__ (celery_app,)在__init__.py中导入的原因很关键它保证 Django 进程启动时 Celery 实例一定被创建后续shared_task装饰器下文详解才能正确解析到当前项目的 app。这种项目模块 任务模块分离的布局适合较大的项目对于简单项目也可以像 First steps with Celery 教程那样用单个模块同时定义 app 与任务。逐行解析celery.py 到底做了什么1. 设置 DJANGO_SETTINGS_MODULE 环境变量os.environ.setdefault(DJANGO_SETTINGS_MODULE, proj.settings)这一行不是必需的但它让你在运行celery命令行程序时无需反复手动传入 settings 模块。它必须出现在创建 app 实例之前。2. 创建 Celery 实例app Celery(proj)proj是应用的主名称同时也是消息默认的前缀。一个项目通常只需一个 Celery 实例即使技术上允许创建多个在 Django 场景下也几乎没有必要。3. 把 Django settings 作为 Celery 配置源app.config_from_object(django.conf:settings, namespaceCELERY)这行代码把 Django 的 settings 对象挂接为 Celery 的配置来源意味着你无需维护多套配置文件直接在 Django 的settings.py里配置 Celery 即可当然两者分离也是允许的。这里有两个要点使用字符串django.conf:settings而非直接传对象文档明确说明用字符串更好因为 worker 进程不需要序列化配置对象。大写命名空间CELERY所有 Celery 配置项必须以大写形式、以CELERY_前缀出现在 Django settings 中。命名空间是可选的但官方强烈推荐以避免与其它 Django settings 命名冲突。4. CELERY_ 前缀的配置映射规则命名空间的含义是Celery 配置项的小写名称如task_always_eager、broker_url、worker_concurrency在 Django settings 中要转换为大写并加CELERY_前缀Celery 配置项小写Django settings 中的写法task_always_eagerCELERY_TASK_ALWAYS_EAGERbroker_urlCELERY_BROKER_URLworker_concurrencyCELERY_WORKER_CONCURRENCYresult_backendCELERY_RESULT_BACKENDtimezoneCELERY_TIMEZONE小写配置名与 DjangoCELERY_前缀 settings 的完整映射关系可参见 User Guide — Configuration 中的 Django namespace 章节。例如一个 Django 项目的settings.py可能包含如下片段与 examples/django/proj/settings.py 的写法一致# settings.py # Celery Configuration Options CELERY_TIMEZONE Australia/Tasmania CELERY_TASK_TRACK_STARTED True CELERY_TASK_TIME_LIMIT 30 * 60示例项目还给出了更完整的一套 broker / 序列化配置examples/django/proj/settings.pyCELERY_BROKER_URL amqp://guest:guestlocalhost # 仅当 broker 处于可信环境时才把 pickle 加入此列表参见 userguide/security.html CELERY_ACCEPT_CONTENT [json] CELERY_RESULT_BACKEND dbsqlite:///results.sqlite CELERY_TASK_SERIALIZER json5. 自动发现任务模块app.autodiscover_tasks()这行代码会让 Celery 自动从所有已安装的 Django app 中发现任务遵循tasks.py约定- app1/ - tasks.py - models.py - app2/ - tasks.py - models.py这样你就不必手动把各个任务模块加入CELERY_IMPORTS配置。从源码看app.autodiscover_tasks()celery/app/base.py在未传packages参数时会委托给 Django fixup处理DjangoFixup.autodiscover_taskscelery/fixups/django.py通过 Django 的 app 注册表枚举所有get_app_configs()并返回其包名随后加载器celery/loaders/base.py为每个包尝试导入相关名.tasks模块。这是一次配置、全部发现机制得以生效的底层原理。6. debug_taskbindTrue 的用法示例app.task(bindTrue, ignore_resultTrue) def debug_task(self): print(fRequest: {self.request!r})debug_task是一个用于调试的示例任务它打印自身请求信息。其中bindTrue是 Celery 3.1 引入的任务选项用于让任务轻松引用当前任务实例此处为self.requestignore_resultTrue表示不保存任务结果。第二步用 shared_task 编写可复用任务真实项目中你的任务很可能放在可复用的 app里。可复用 app 不能依赖具体项目本身因此也不能直接 import 项目的 Celery 实例。shared_task装饰器正是为此而生——它允许你在没有任何具体 app 实例的情况下创建任务任务会从当前 app 的任务注册表中动态解析。仓库中 examples/django/demoapp/tasks.py 展示了完整的任务集合from demoapp.models import Widget from celery import shared_task shared_task def add(x, y): return x y shared_task def mul(x, y): return x * y shared_task def xsum(numbers): return sum(numbers) shared_task def count_widgets(): return Widget.objects.count() shared_task def rename_widget(widget_id, name): w Widget.objects.get(idwidget_id) w.name name w.save() shared_task( bindTrue, autoretry_for(Exception,), retry_kwargs{max_retries: 2, countdown: 10 * 60}, # 最多重试 2 次间隔 10 分钟 ) def error_task(self): raise Exception(Test error) shared_task( bindTrue, autoretry_for(Exception,), retry_backoff5, # 退避因子秒第一次 5s、第二次 10s、第三次 20s…… retry_jitterFalse, # 设为 False 禁用随机抖动使用精确值5s、10s、20s retry_kwargs{max_retries: 3}, ) def error_backoff_test(self): raise Exception(Test error)可以看到shared_task不仅支持纯函数式任务add、mul、xsum也能配合 Django ORM 操作数据库count_widgets、rename_widget还能通过bindTrueautoretry_for/retry_backoff等选项实现带指数退避的自动重试。在底层shared_taskcelery/app/init.py返回一个Proxy代理对象任务总是从当前 app的任务注册表中取因此同一个shared_task定义可以服务于不同项目、不同 app 实例这正是可复用 app 场景所需要的。常见陷阱任务触发与数据库事务问题任务可能在事务提交前运行Django 集成中最常见的坑是在视图中立即触发任务却不等待数据库事务提交。此时 Celery worker 可能在变更持久化到数据库之前就执行任务导致任务读不到数据。例如# views.py def create_user(request): # 注意简化示例实际应使用表单校验输入 user User.objects.create(usernamerequest.POST[username]) send_email.delay(user.pk) return HttpResponse(User created) # task.py shared_task def send_email(user_pk): user User.objects.get(pkuser_pk) # send email ...在这个例子中send_email任务可能在视图提交事务前就启动因此任务可能找不到刚创建的用户。方案一Django 的 transaction.on_commit 钩子常见的解决方案是使用 Django 的transaction.on_commit钩子在事务提交后再触发任务- send_email.delay(user.pk) transaction.on_commit(lambda: send_email.delay(user.pk))方案二Celery 5.4 的 delay_on_commit 快捷方法由于提交后再触发是如此常见的模式Celery 5.4 起提供了一个便捷 API使用celery.contrib.django.task.DjangoTask。不要调用delay()而是调用delay_on_commit()- send_email.delay(user.pk) send_email.delay_on_commit(user.pk)该方法替你完成了把调用包装进on_commit钩子的工作。对应源码在 celery/contrib/django/task.pyimport functools from django.db import transaction from celery.app.task import Task class DjangoTask(Task): Extend the base :class:~celery.app.task.Task for Django. Provide a nicer API to trigger tasks at the end of the DB transaction. def delay_on_commit(self, *args, **kwargs) - None: Call :meth:~celery.app.task.Task.delay with Djangos on_commit(). transaction.on_commit(functools.partial(self.delay, *args, **kwargs)) def apply_async_on_commit(self, *args, **kwargs) - None: Call :meth:~celery.app.task.Task.apply_async with Djangos on_commit(). transaction.on_commit(functools.partial(self.apply_async, *args, **kwargs))需要注意delay_on_commit与delay的关键差异delay_on_commit不会向调用方返回任务 ID。任务在调用该方法时并不会被发送到 broker只有 Django 事务结束时才会真正发出如果你需要任务 ID请继续使用delay()在极少数需要不等待事务、立即触发的场景下原有的delay()API 仍然可用。如果你完全按照本文的 setup 步骤操作DjangoTask会被自动使用从 celery/fixups/django.py 可以看到Django fixup 安装时会把app.task_cls设置为celery.contrib.django.task:DjangoTask。但如果你使用了自定义任务基类参见 User Guide — Tasks 的自定义任务类则需要显式继承DjangoTask而不是celery.app.task.Task才能获得此行为from celery.contrib.django.task import DjangoTask class MyTask(DjangoTask): ...单元测试对上述行为有明确覆盖t/unit/contrib/django/test_task.py中的test_delay_on_commit断言delay_on_commit()返回None并验证调用被包装进transaction.on_commit。Django 连接池支持Django 5.1从Django 5.1起Django 内置了数据库连接池支持。如果你在 Django 的DATABASES设置中启用了连接池Celery 会自动在 worker 进程中处理连接池的关闭具体通过数据库后端的close_pool方法实现——因为数据库连接无法跨进程共享。这一机制在 celery/fixups/django.py 中有完整实现DjangoWorkerFixup._close_database在遍历连接时会检查该连接的DATABASES[alias][OPTIONS][pool]是否启用若启用且 worker 使用 prefork 池则调用conn.close_pool()。worker 进程在task_prerun/task_postruncelery/fixups/django.py以及进程初始化时on_worker_process_init都会关闭数据库与缓存连接避免跨进程复用导致的连接损坏。扩展django-celery-results——用 Django ORM/Cache 作为结果后端django-celery-results扩展提供了基于Django ORM或Django Cache 框架的结果后端。接入步骤1. 安装库$ pip install django-celery-results2. 在settings.py的INSTALLED_APPS中加入django_celery_resultsINSTALLED_APPS ( ..., django_celery_results, )注意模块名中是下划线而非连字符django_celery_results。3. 执行数据库迁移创建 Celery 数据表$ python manage.py migrate django_celery_results4. 配置 Celery 使用该后端假设仍通过 Django 的settings.py配置 Celery使用 ORM 后端CELERY_RESULT_BACKEND django-db使用缓存后端时可以指定CACHES设置中定义的某个缓存CELERY_RESULT_BACKEND django-cache # 从 CACHES 设置中选择使用哪个缓存 CELERY_CACHE_BACKEND default # django 设置 CACHES { default: { BACKEND: django.core.cache.backends.db.DatabaseCache, LOCATION: my_cache_table, } }更多的结果后端配置选项参见 User Guide — Configuration 的结果后端章节。扩展django-celery-beat——数据库驱动的周期任务django-celery-beat扩展提供数据库存储的周期任务与Admin 管理界面适合在 Django Admin 中动态管理定时任务。详细用法参见 User Guide — Periodic Tasks 的自定义调度器章节。启动 worker 进程在生产环境中你应当把 worker 作为后台守护进程运行参见 Daemonizing但在测试和开发阶段可以使用celery worker管理命令直接启动一个 worker 实例就像使用 Django 的manage.py runserver一样$ celery -A proj worker -l INFO其中-A proj指定了包含 Celery 实例的模块即我们在proj/proj/celery.py中定义的app。查看所有可用的命令行选项$ celery --help源码视角Django fixup 是如何自动生效的当我们设置了DJANGO_SETTINGS_MODULE环境变量后Celery 会自动安装 Django 集成celery/fixups/django.py 的fixup()函数只要环境变量存在且未使用自定义加载器Celery 就会尝试导入 Django、校验版本并安装DjangoFixup。该 fixup 主要做几件事把当前工作目录加入sys.path保证能 import 到项目模块将 app 的时间now指向django.utils.timezone.now使任务时间与 Django 的时区设置保持一致在未使用自定义任务类时将app.task_cls切换为DjangoTask这正是delay_on_commit自动可用的原因连接import_modules信号在导入任务模块前调用django.setup()并执行 Django 系统检查可通过环境变量CELERY_SKIP_CHECKS跳过连接worker_init等信号安装 worker 进程级别的数据库/缓存连接管理DjangoWorkerFixup。下一步学习完成本文的集成后你可以继续阅读 Next Steps 教程随后深入 User Guide 学习调用方式、Canvas 工作流、路由、监控等进阶主题。完整的 Django 示例项目源码含proj工程与demoapp应用就在本仓库的 examples/django 目录下可作为你动手实践的直接参考。【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网