django-simpletask5
本项目由 opencode + deepseek-v4-flash 生成
一个轻量级的 Django 异步任务执行框架,提供声明式的任务模型、信号驱动的自动发布、Worker 进程异步执行,以及内置的 Cron 定时调度。
特性
- 声明式任务模型 — 继承
Task模型即可定义任务,自动处理创建/更新/删除事件 - 信号驱动 — Django 信号自动拦截模型变更,发布
TaskExecution到消息队列 - 自定义事件 — 通过
task.trigger('event_name')触发任意事件 - 灵活的执行器映射 — 不同事件可绑定不同的执行器类
- 队列路由 — 不同事件可路由到不同优先级队列
- 重试与超时 — 失败自动重试(指数退避),支持执行超时与 Pickup 超时两种检测
- 拒绝回退 — 执行器可通过
RejectException拒绝当前任务,消息自动重新入队等待下次调度 - 加密字段 — 敏感数据自动加密存储
- Cron 调度 — 内置 crontab 守护进程,支持代码注册与数据库覆盖
- 归档统计 — 已完成执行记录自动归档为加密 JSONL,并生成日统计
- Worker 注册中心 — 基于 Redis 的 Worker 心跳与状态追踪
- Django Admin 集成 — 完整的后台管理界面
依赖
- Python >= 3.8
- Django >= 3.2(兼容至 5.2.x)
- Kombu >= 5.3.0(消息队列,支持 RabbitMQ/Redis/内存)
- Redis >= 4.0.0(分布式锁、Worker 注册中心)
- 详见
pyproject.toml
安装
pip install django-simpletask5
django-simpletask5 使用 django-app-requires 自动管理依赖的 app(如 django_safe_fields),无需手动添加到 INSTALLED_APPS。只需在 settings.py 中调用一次 patch:
from django_app_requires import patch_all as django_app_requires_patch_all
django_app_requires_patch_all()
然后在 INSTALLED_APPS 中添加:
INSTALLED_APPS = [
...
'django_simpletask5',
]
配置分布式锁(最小可用配置,完整配置项见下方「配置说明」):
DJANGO_SIMPLETASK_LOCK_CONFIG = {
'global_lock_engine_class': 'globallock.redis_global_lock.RedisGlobalLock',
'global_lock_engine_options': {
'host': '127.0.0.1',
'port': 6379,
'db': 0,
},
}
运行迁移:
python manage.py migrate django_simpletask5
配置说明
以下给出面向项目集成的完整 settings.py 配置示例。除 django-simpletask5 相关配置外,其余 Django 默认配置项(DATABASES、MIDDLEWARE、TEMPLATES、ROOT_URLCONF 等)按项目自身需要配置,此处以 ... 省略。配置原则逐项见后文「配置项与配置原则」表。
import os
from pathlib import Path
BASE_DIR = Path(__file__).resolve().parent.parent
SECRET_KEY = os.environ.get('DJANGO_SECRET_KEY', 'change-me-in-production')
DEBUG = os.environ.get('DJANGO_DEBUG', 'False') == 'True'
ALLOWED_HOSTS = ['*']
# ... 数据库(DATABASES)、中间件、模板、URL 等 Django 默认配置 ...
INSTALLED_APPS = [
# django-simpletask5 的依赖 app(django_safe_fields 等)
# 由 django-app-requires 自动补全,无需手动添加
'django_simpletask5',
# ... 你的业务 app ...
]
# ── 消息队列(MQ)────────────────────────────────────────────
# 生产使用 RabbitMQ 或 Redis;开发调试可用内存队列
DJANGO_SIMPLETASK_BROKER_URL = 'amqp://guest:guest@127.0.0.1:5672//'
# 若不设 BROKER_URL,则按以下参数组合:
# DJANGO_SIMPLETASK_MESSAGE_QUEUE_TYPE = 'rabbitmq' # 'memory' | 'rabbitmq' | 'redis'
# DJANGO_SIMPLETASK_RABBITMQ_HOST = '127.0.0.1'
# DJANGO_SIMPLETASK_RABBITMQ_PORT = 5672
# DJANGO_SIMPLETASK_RABBITMQ_USER = 'guest'
# DJANGO_SIMPLETASK_RABBITMQ_PASSWORD = 'guest'
# DJANGO_SIMPLETASK_REDIS_MQ_HOST = '127.0.0.1' # MESSAGE_QUEUE_TYPE='redis' 时
# DJANGO_SIMPLETASK_REDIS_MQ_PORT = 6379
# DJANGO_SIMPLETASK_REDIS_MQ_DB = 0
# DJANGO_SIMPLETASK_REDIS_MQ_PASSWORD = None # Redis MQ 需认证时设置
# ── 分布式锁与 Worker 注册中心(Redis)────────────────────────
# 默认锁引擎 DjangoRedisGlobalLock 复用 CACHES['default'],需为 Redis 缓存
CACHES = {
'default': {
'BACKEND': 'django_redis.cache.RedisCache',
'LOCATION': 'redis://127.0.0.1:6379/1',
'OPTIONS': {'CLIENT_CLASS': 'django_redis.client.DefaultClient'},
},
}
# 如需独立连接参数(锁与 Worker 注册中心共用同一 host/port/db/password):
# DJANGO_SIMPLETASK_LOCK_CONFIG = {
# 'global_lock_engine_class': 'globallock.redis_global_lock.RedisGlobalLock',
# 'global_lock_engine_options': {'host': '127.0.0.1', 'port': 6379, 'db': 0,
# 'password': '<redis 密码,无需认证可省略>'},
# }
# 注:globallock>=0.2.0 已取消按 timeout 自动释放的概念。锁的有效期由
# lockman/lock 作用域自行管理,作用域未结束前锁持续有效(内置看门狗自动续期),
# 锁在 `with` 块退出或进程消亡时自动释放。已移除 DJANGO_SIMPLETASK_LOCK_TIMEOUT
# 与 DJANGO_SIMPLETASK_LOCK_RENEW_INTERVAL,请勿再配置这两项。
# ── Key 前缀与队列路由 ───────────────────────────────────────
DJANGO_SIMPLETASK_KEY_PREFIX = 'myproject' # Redis key 前缀,默认取 ROOT_URLCONF 首段
DJANGO_SIMPLETASK_DEFAULT_QUEUE = 'myproject:django-simpletask5:queue:default'
DJANGO_SIMPLETASK_HIGH_PRIORITY_QUEUE = 'myproject:django-simpletask5:queue:high_priority'
# ── 执行与重试 ───────────────────────────────────────────────
DJANGO_SIMPLETASK_DEFAULT_MAX_RETRIES = 3
DJANGO_SIMPLETASK_DEFAULT_TIMEOUT_SECONDS = 600 # 执行超时(秒)
DJANGO_SIMPLETASK_DEFAULT_PICKUP_TIMEOUT_SECONDS = 3600 # Pickup 超时(秒)
# ── Worker / 日志运维参数 ────────────────────────────────────
DJANGO_SIMPLETASK_WORKER_HEARTBEAT_INTERVAL = 60 # Worker 心跳间隔(秒)
DJANGO_SIMPLETASK_SHUTDOWN_TIMEOUT = 30 # 优雅停机等待(秒)
DJANGO_SIMPLETASK_LOG_INTERVAL = 300 # 存活/节流日志间隔(秒)
# ── 归档与统计 ──────────────────────────────────────────────
DJANGO_SIMPLETASK_ARCHIVE_PATH = os.path.join(BASE_DIR, 'django_simpletask5_archives')
DJANGO_SIMPLETASK_ARCHIVE_SALT = 'change-me-archive-salt' # 归档加密盐值
DJANGO_SIMPLETASK_ARCHIVE_RETENTION_DAYS = 7
# ── Cron ────────────────────────────────────────────────────
DJANGO_SIMPLETASK_CRONJOB_AUTO_SYNC = True # 自动同步代码注册的 Cron
# ── 安全:脚本执行器与字段加密 ───────────────────────────────
DJANGO_SIMPLETASK_ENABLE_SCRIPT_EXECUTORS = False # 脚本执行器开关(生产保持关闭)
DJANGO_SIMPLETASK_SCRIPT_WHITELIST = [] # 脚本白名单
DJANGO_SIMPLETASK_FIELD_CIPHER_CLASS = None # 自定义加密器类
DJANGO_SIMPLETASK_RESULT_PASSWORD = os.environ.get('DJANGO_SIMPLETASK_RESULT_PASSWORD')
DJANGO_SIMPLETASK_ERROR_MSG_PASSWORD = os.environ.get('DJANGO_SIMPLETASK_ERROR_MSG_PASSWORD')
DJANGO_SIMPLETASK_CONTEXT_PASSWORD = os.environ.get('DJANGO_SIMPLETASK_CONTEXT_PASSWORD')
# ── 可选:Admin Dashboard ────────────────────────────────────
# DJANGO_ADMIN_DASHBOARDS = {
# 'admin:app_list': {'django_simpletask5': 'django_simpletask5.dashboards.TaskDashboard'},
# }
# 必须最后调用:自动补全 django-simpletask5 依赖的 app
from django_app_requires import patch_all as django_app_requires_patch_all
django_app_requires_patch_all()
配置项与配置原则
| 配置项 | 默认值 | 配置原则 |
|---|---|---|
DJANGO_SIMPLETASK_BROKER_URL |
None |
优先级最高,设置后忽略下方 MQ 类型/连接参数。生产建议统一使用完整 URL,便于复用云厂商连接串与凭据管理 |
DJANGO_SIMPLETASK_MESSAGE_QUEUE_TYPE |
'memory' |
memory 仅用于开发调试;生产必须改为 rabbitmq 或 redis,否则消息不持久化、跨进程 Worker 无法消费 |
DJANGO_SIMPLETASK_RABBITMQ_HOST/PORT/USER/PASSWORD |
127.0.0.1 / 5672 / guest / guest |
类型为 rabbitmq 且未设 BROKER_URL 时生效。生产务必更换默认 guest 凭据 |
DJANGO_SIMPLETASK_REDIS_MQ_HOST/PORT/DB/PASSWORD |
127.0.0.1 / 6379 / 0 / None |
类型为 redis 时生效。建议与锁/缓存使用不同 db,避免 key 冲突与互相阻塞。Redis 需认证时设置 PASSWORD |
CACHES['default'] |
— | 默认锁引擎 DjangoRedisGlobalLock 依赖它,必须配置为 Redis 缓存 |
DJANGO_SIMPLETASK_LOCK_CONFIG |
None(默认复用缓存锁) |
需要独立连接参数时显式配置(global_lock_engine_options 支持 host/port/db/password)。注意:Worker 注册中心(心跳/状态)也读取该配置的 host/port/db/password,未配置时固定连 127.0.0.1:6379/0。globallock>=0.2.0 起锁有效期由 lockman/lock 作用域管理(内置看门狗自动续期),不再需要配置 LOCK_TIMEOUT / LOCK_RENEW_INTERVAL |
DJANGO_SIMPLETASK_KEY_PREFIX |
自动取 ROOT_URLCONF 首段 |
生产建议显式指定,多项目共用 Redis 时避免 key 冲突 |
DJANGO_SIMPLETASK_DEFAULT_QUEUE |
{prefix}:...:queue:default |
默认队列名。多环境(dev/staging/prod)建议加环境后缀隔离 |
DJANGO_SIMPLETASK_HIGH_PRIORITY_QUEUE |
{prefix}:...:queue:high_priority |
高优先级队列,需单独启动 Worker 消费 |
DJANGO_SIMPLETASK_DEFAULT_MAX_RETRIES |
3 |
默认最大重试次数。执行器幂等性好的任务可加大 |
DJANGO_SIMPLETASK_DEFAULT_TIMEOUT_SECONDS |
600 |
执行超时,应大于任务实际最长执行时间,否则任务被误判超时 |
DJANGO_SIMPLETASK_DEFAULT_PICKUP_TIMEOUT_SECONDS |
3600 |
Pickup 超时(等待 Worker 消费的上限),通常大于执行超时,因为排队等待时间不可控 |
DJANGO_SIMPLETASK_WORKER_HEARTBEAT_INTERVAL |
60 |
Worker 心跳间隔,应小于 Worker 注册 TTL(300 秒),否则注册状态被误清理 |
DJANGO_SIMPLETASK_SHUTDOWN_TIMEOUT |
30 |
优雅停机等待时长,应大于单任务最长执行时长,否则停机会中断进行中的任务 |
DJANGO_SIMPLETASK_LOG_INTERVAL |
300 |
存活日志与异常节流日志的输出间隔。日志量大时可调大 |
DJANGO_SIMPLETASK_ARCHIVE_PATH |
django_simpletask5_archives |
归档目录。生产建议指向持久化/对象存储挂载目录 |
DJANGO_SIMPLETASK_ARCHIVE_SALT |
django-simpletask5-archive |
归档加密盐值,生产必须修改(结合 SECRET_KEY 派生 AES 密钥) |
DJANGO_SIMPLETASK_ARCHIVE_RETENTION_DAYS |
7 |
归档保留天数,按合规要求调整 |
DJANGO_SIMPLETASK_CRONJOB_AUTO_SYNC |
True |
自动把代码注册的 Cron 同步到数据库;只想手工管理时设为 False |
DJANGO_SIMPLETASK_ENABLE_SCRIPT_EXECUTORS |
False |
Python/Shell 脚本执行器开关,生产保持 False;开启还需授予对应权限 |
DJANGO_SIMPLETASK_SCRIPT_WHITELIST |
[] |
脚本执行白名单(空列表表示不限制) |
DJANGO_SIMPLETASK_FIELD_CIPHER_CLASS |
None |
自定义字段加密器类(默认 django-safe-fields 内置加密) |
DJANGO_SIMPLETASK_RESULT_PASSWORD |
None |
result 字段加密密码,建议通过环境变量注入,勿硬编码 |
DJANGO_SIMPLETASK_ERROR_MSG_PASSWORD |
None |
error_message 字段加密密码,同上 |
DJANGO_SIMPLETASK_CONTEXT_PASSWORD |
None |
context 字段加密密码,同上 |
注意:
DJANGO_SIMPLETASK_ENABLE_SCRIPT_EXECUTORS设为True后,还需为用户/组授予django_simpletask5 | Cron job | Can use script executors权限,非超级管理员无法创建 Python/Shell 脚本执行器类型的定时任务。
快速开始(从零到一集成)
本框架的集成分三步:定义任务模型 → 编写执行器 → 启动 Worker。完整可运行的参考实现见仓库中的 django_simpletask5_example 应用,包含 OrderTask 模型、创建/更新/删除/退款四个执行器及内置 Cron 任务。
1. 定义任务模型
继承 django_simpletask5.models.Task 即可获得信号驱动的异步能力。模型的保存(新增/更新)与删除会被自动拦截并发布消息,无需手写信号:
from django.db import models
from django_simpletask5.models import Task
class OrderTask(Task):
order_id = models.CharField(max_length=64, unique=True)
customer_name = models.CharField(max_length=128)
amount = models.DecimalField(max_digits=10, decimal_places=2)
status = models.CharField(max_length=32, default='pending')
# 事件 → 执行器映射(详见「事件处理机制」)
executor_class = {
'create': 'myapp.executors.OrderCreateExecutor',
'update': 'myapp.executors.OrderUpdateExecutor',
'delete': 'myapp.executors.OrderDeleteExecutor',
'refund': 'myapp.executors.OrderRefundExecutor', # 自定义事件,由 request_refund 通过 self.trigger('refund') 手动触发
}
# 事件 → 队列路由(未映射的事件自动落到 DJANGO_SIMPLETASK_DEFAULT_QUEUE)
simpletask_queue = {
'create': 'myproject:django-simpletask5:queue:high_priority',
'update': 'myproject:django-simpletask5:queue:default',
}
# 仅当这些字段发生变化时才触发 update 事件(详见「事件处理机制」)
trigger_update_fields = ['customer_name', 'amount', 'status']
def request_refund(self, amount):
return self.trigger('refund', extra_context={'refund_amount': str(amount)})
要点:
- 主键为 Django 默认自增
id;内置task_id(UUID,unique=True)作为业务唯一标识,不是主键 - 保存 / 删除 / 批量创建自动触发事件,详见「事件处理机制」
- 事件上下文、update 字段判定、delete 拦截等行为的完整说明见对应章节
2. 编写执行器
执行器继承 BaseExecutor,实现 execute() 方法即可:
# myapp/executors.py
from django_simpletask5.executors.base import BaseExecutor
from django_simpletask5.models import TaskExecution
class OrderCreateExecutor(BaseExecutor):
timeout_seconds = 300 # 执行超时 5 分钟(可选)
pickup_timeout_seconds = 1800 # Pickup 超时 30 分钟(可选)
def execute(self, execution: TaskExecution, task=None) -> str | None:
context = execution.get_context_dict() # 事件上下文
order = context['data'] # create 事件的模型字段快照
# ... 你的业务逻辑(调用外部系统、发消息等)...
return 'ok' # 返回值存入 execution.result
执行器接口、批量任务、异常重试、拒绝回退等完整说明见「应用开发:执行器」。
3. 启动 Worker
python manage.py django_simpletask_executor --workers 4
- 默认监听
default队列;high_priority队列需单独启动 Worker 处理 - 可用
--service myapp.executors.OrderCreateExecutor只处理指定执行器 - 可用
--debug输出完整异常堆栈与 DEBUG 日志 - Worker 线程由
WorkerSupervisor守护:异常退出自动重启、保持目标 Worker 数量;Redis/MQ 断连自动重连
4. 启动 Cron 调度(可选)
python manage.py django_simpletask_crontab
5. 触发自定义事件
在模型方法或业务代码中调用 trigger() 即可发布任意事件:
order = OrderTask.objects.get(order_id='ORD-001')
order.trigger('refund', extra_context={'refund_amount': '50.00'})
6. 验证
启动 Worker 后新增一条 OrderTask,观察数据流转:
OrderTask.objects.create(order_id='ORD-001', customer_name='Test', amount=Decimal('100.00'))
- 数据库
TaskExecution表新增一条create事件、pending状态的记录 - Worker 消费后状态变为
running,执行完成变为success,result保存执行器返回值 - Django Admin 的 Dashboard 可查看执行统计与 Worker 状态
事件处理机制
内置事件:create / update / delete
框架通过 Django 信号自动监听所有 Task 子类的保存与删除,无需手动发布:
| 事件 | 触发时机 | 执行器可见的上下文 |
|---|---|---|
create |
新增实例(save() / bulk_create()) |
task_id、trigger_event、data(全字段快照) |
update |
字段值发生变化后保存 | task_id、trigger_event、data、old_data、new_data、changed_fields |
delete |
调用 instance.delete() |
task_id、trigger_event、data |
执行器通过 execution.get_context_dict() 获取以上上下文。
update 事件的字段变更判定
update 事件与保存动作的关系:框架每次保存已存在的实例时,通过 pre_save/post_save 信号对比旧值快照与新值,计算发生变化的字段列表(changed_fields)。默认只要任一字段变化就发布 update 事件——这往往会产生大量非预期消息。例如某次保存只改了一个内部计数或审计字段,下游系统却会收到一次同步通知。
为此可用类属性控制「哪些字段变化才值得触发 update 事件」:
class OrderTask(Task):
# 白名单:仅当这些字段之一发生变化时才发布 update 事件
trigger_update_fields = ['customer_name', 'amount', 'status']
# 或黑名单:这些字段变化不发布 update 事件,其余字段变化照常发布
# trigger_update_blacklist = ['sync_count', 'updated_by']
| 配置方式 | 行为 | 适用场景 |
|---|---|---|
| 都不配置(默认) | 任何字段变化都触发 update 事件 | 简单模型 |
trigger_update_fields |
仅当白名单内的字段变化时触发 | 只在关键业务字段变更时同步下游(推荐) |
trigger_update_blacklist |
白名单外的字段变化触发;黑名单内字段变化不触发 | 忽略噪声字段(内部计数、审计字段等) |
白名单与黑名单二选一,同时配置时以白名单为准。无论何种方式,值未实际变化(如调用 save() 但所有字段值相同)的保存都不会触发 update 事件。changed_fields 中即包含本次实际发生变化的字段名列表,执行器可据此做增量处理。
delete 事件的拦截与确认删除
Task.delete() 被框架拦截,不会立即删除数据库记录,而是发布 delete 事件消息。执行器处理完外部系统清理后,调用 execution.confirm_delete() 完成真实删除:
class OrderDeleteExecutor(BaseExecutor):
def execute(self, execution: TaskExecution, task=None) -> str | None:
# ... 先清理外部系统关联数据 ...
execution.confirm_delete() # 真正删除对应 Task 记录
return 'deleted'
这样实现删除操作与外部系统清理的解耦,且支持失败重试。
事件 → 执行器映射
每个事件消息最终要由一个执行器处理。executor_class 声明「事件名 → 执行器类路径」的分派表,可以是字符串(所有事件共用同一个执行器)或字典(按事件分派):
executor_class = {
'create': 'myapp.executors.OrderCreateExecutor',
'update': 'myapp.executors.OrderUpdateExecutor',
'delete': 'myapp.executors.OrderDeleteExecutor',
'refund': 'myapp.executors.OrderRefundExecutor', # 自定义事件
'default': 'myapp.executors.OrderDefaultExecutor', # 兜底执行器(可选)
}
执行器解析顺序:命中事件名(如 refund 映射到 OrderRefundExecutor)→ default → create → 自动派生 app_label.tasks.{ClassName}{Event.capitalize()}Executor(以上都未命中时的兜底路径)。
注意:
executor_class只负责「事件 → 由谁执行」的映射,与事件由谁触发无关。内置事件(create/update/delete)由信号自动触发;自定义事件(如refund)必须由业务代码调用trigger('refund')手动触发,详见下方「自定义事件」。
事件 → 队列路由
每个事件消息发布到哪个队列,由 simpletask_queue 决定。可以是字符串(所有事件同一队列)或字典(按事件路由):
simpletask_queue = {
'create': 'myproject:django-simpletask5:queue:high_priority',
'update': 'myproject:django-simpletask5:queue:default',
}
队列解析顺序:命中事件名 → default → 全局默认队列 DJANGO_SIMPLETASK_DEFAULT_QUEUE。
以上面只配置 create/update 两个键为例,各事件的落队结果:
| 事件 | 解析结果 |
|---|---|
create |
命中键 → 高优先级队列 |
update |
命中键 → 默认队列 |
delete(内置) |
未命中、无 default 键 → 全局默认队列 DJANGO_SIMPLETASK_DEFAULT_QUEUE |
refund(自定义) |
同上 → 全局默认队列 |
即:只映射了部分事件时,其余未映射的事件自动落到 DJANGO_SIMPLETASK_DEFAULT_QUEUE。如需统一兜底到某个队列,可显式添加 'default' 键。
自定义事件
与信号自动触发的内置事件不同,自定义事件(如 refund)不会自动触发,需要业务代码调用 Task.trigger(event, extra_context=None) 手动发布:
order = OrderTask.objects.get(order_id='ORD-001')
order.trigger('refund', extra_context={'refund_amount': '50.00'})
# 触发后流程:
# trigger('refund') → 创建 trigger_event='refund' 的 TaskExecution
# → 按 simpletask_queue 路由发布消息 → Worker 按 executor_class['refund']
# 加载 OrderRefundExecutor 执行
- 事件名任意:
trigger()的第一个参数即事件名,需与executor_class/simpletask_queue字典中的键保持一致,否则走default或兜底路径 - 执行器:按「事件 → 执行器映射」解析
- 队列:按「事件 → 队列路由」解析
- 上下文:
extra_context会合并进执行器可见的事件上下文,执行器侧读取context['refund_amount']:
class OrderRefundExecutor(BaseExecutor):
def execute(self, execution, task=None):
context = execution.get_context_dict()
refund_amount = context['refund_amount'] # == '50.00'
...
抑制异步消息
临时禁止某个实例发布消息(_skip_publish 为临时实例属性,不会持久化到数据库),四种使用方式:
order._skip_publish = True; order.save() # 方式一:实例属性
order.skip_publish().save() # 方式二:链式方法
order.save(skip_publish=True) # 方式三:save 参数
OrderTask.objects.create(..., skip_publish=True) # 方式四:构造参数(推荐)
不会创建 TaskExecution,也不会发送异步消息。
应用开发:执行器
BaseExecutor 接口
| 成员 | 说明 |
|---|---|
execute(execution, task=None) |
必须实现。执行任务,返回值存入 execution.result |
get_tasks(execution) |
返回可迭代对象,默认 yield None(单任务) |
get_tasks_count(execution) |
任务总数(用于进度展示),默认 1 |
reject(execution, reason='') |
拒绝任务:状态复位 pending、清空 started_at、刷新 expire_time |
timeout_seconds / pickup_timeout_seconds |
类属性,覆盖该执行器的超时配置 |
get_parameter_schema() |
声明执行器参数,用于 Admin 参数配置界面(默认空) |
execute 签名与返回值
def execute(self, execution: TaskExecution, task=None) -> str | None:
context = execution.get_context_dict() # 事件上下文(见「事件处理机制」)
# ... 业务逻辑 ...
return 'result'
- 事件上下文:
execution.get_context_dict()返回反序列化后的 JSON 上下文 - 返回值:单个任务时直接返回(字符串或 JSON 字符串);返回
None时result保持为空
批量任务
重写 get_tasks() 与 get_tasks_count(),一个执行器即可处理多条子任务:
class BatchSyncExecutor(BaseExecutor):
def get_tasks_count(self, execution):
return len(self._ids(execution))
def get_tasks(self, execution):
for _id in self._ids(execution):
yield _id
def _ids(self, execution):
return execution.get_context_dict().get('ids', [])
def execute(self, execution, task=None):
self._sync(task) # task 为 get_tasks() 产出的单个子任务
return 'done'
批量执行时 execution.result 聚合为 {"results": [...], "total": N},每完成一个子任务更新 done_tasks_count 便于进度展示。
异常与状态流转
执行器抛出的异常由框架统一处理:
| 情况 | 结果 |
|---|---|
| 无异常 | 状态 success,result 保存返回值 |
RejectException |
拒绝回退:状态复位 pending,消息重新入队,不消耗重试次数 |
其他异常且 retry_count < max_retries |
状态 retry,retry_count + 1,消息重新入队 |
| 其他异常且重试次数已用完 | 状态 failed,error_message 记录异常信息 |
拒绝回退
执行器遇到临时条件不满足(如资源上限),抛出 RejectException 拒绝本次执行:
from django_simpletask5.executors.exceptions import RejectException
class MyExecutor(BaseExecutor):
def execute(self, execution: TaskExecution, task=None) -> str | None:
if not self._can_proceed():
raise RejectException('Resource limit reached, try again later')
框架捕获后:
- 调用
message.reject(requeue=True)将消息放回队列尾部,下次重新调度 TaskExecution状态恢复为pending,started_at清空,等待下次 Worker 消费- 不消耗重试次数
- 刷新
expire_time(使用 Pickup 超时),防止被StatusCheckExecutor误回收
超时配置
class OrderCreateExecutor(BaseExecutor):
timeout_seconds = 300 # 执行超时 5 分钟
pickup_timeout_seconds = 1800 # Pickup 超时 30 分钟
未配置时使用全局默认值 DJANGO_SIMPLETASK_DEFAULT_TIMEOUT_SECONDS(600 秒)与 DJANGO_SIMPLETASK_DEFAULT_PICKUP_TIMEOUT_SECONDS(3600 秒)。
模型与生命周期
Task 抽象模型
所有业务任务模型继承 Task:
- 内置
task_id(UUID,unique=True);主键为 Django 默认自增id,无需自定义 - 提供
trigger(event, extra_context)、skip_publish()、get_executor_class_path(event)、get_queue_name(event)等方法 - 保存、删除、批量创建自动触发事件(见「事件处理机制」)
TaskExecution 状态机
每条事件消息对应一条 TaskExecution 记录,状态流转如下:
pending ──(Worker 消费)──▶ running ──(执行成功)──▶ success
│
├─(执行异常且重试次数用完)──▶ failed
├─(执行异常但可重试)────────▶ retry ──(重新入队)──▶ pending
└─(expire_time 超时)────────▶ timeout ──(重试)────▶ retry
pending ──(用户取消)────────▶ canceled(终态)
| 状态 | 含义 |
|---|---|
pending |
已发布,等待 Worker 消费 |
running |
正在执行 |
success |
执行成功 |
failed |
重试次数耗尽后失败 |
retry |
失败但可重试,消息已重新入队 |
timeout |
超过 expire_time 未完成 |
canceled |
被用户主动取消(终态,Worker 不会执行) |
取消任务
已发布但未执行(pending)的任务可通过 Admin 列表页勾选后执行 「Cancel selected pending tasks」 操作批量取消,状态置为 canceled。代码方式:
from django_simpletask5.models import TaskExecution
TaskExecution.objects.filter(
execution_id='<execution_id>', status='pending',
).update(status='canceled')
canceled 为终态:即使消息仍留在 MQ 队列中,Worker 消费时也会直接跳过(不会执行),且不会被归档统计混入 failed。已进入 running 的任务无法取消。
重试与超时
- 重试:
max_retries(默认 3)控制最大重试次数,retry_count累计已重试次数;重试消息重新发布到原队列 - 两种超时:
- 执行超时(
timeout_seconds):进入running后必须在指定时间内完成,防止执行 hang 死 - Pickup 超时(
pickup_timeout_seconds):创建后被 Worker 消费的等待上限,也用于拒绝回退后的重新调度保护
- 执行超时(
- 超时检测与恢复:内置
StatusCheckExecutor(内置 Cron)将超时的pending/running记录标记为timeout;RetryTimeoutExecutor将timeout记录重新入队重试(执行超时消耗重试次数,Pickup 超时不消耗)
其他模型
CronJob— 定时任务注册表,支持代码注册与数据库覆盖(见「Cron 任务」)TaskExecutionArchive— 归档记录,已完成执行归档为加密 JSONL(见配置「归档」)TaskExecutionStat— 按executor_class与日期聚合的日统计
Cron 任务
在代码中定义 Cron 任务
在任意 app 的 cronjobs.py 文件中使用 register_cronjob 注册:
# myapp/cronjobs.py
from django_simpletask5.cronjob_registry import register_cronjob
register_cronjob(
name='health_check',
cron_expression='*/5 * * * *',
executor_class='django_simpletask5.executors.simple_request.SimpleRequestExecutor',
context={
'url': 'https://example.com/health',
'method': 'GET',
'timeout': 10,
},
description='定期健康检查',
)
从代码同步到数据库
注册的 Cron 任务需要同步到数据库才会生效。框架默认在 django_simpletask_crontab 启动时自动同步(可通过 DJANGO_SIMPLETASK_CRONJOB_AUTO_SYNC = False 关闭),也可以手动执行:
python manage.py django_simpletask_sync_cronjobs
同步后,用户可以在 Django Admin 中查看和修改 Cron 任务,被手动修改过的任务不会在后续同步中被覆盖(is_modified_by_user 标记保护)。
注意:cron 表达式按 Django 的
TIME_ZONE设置对应的本地时间解释。USE_TZ=True时仍以本地时间计算触发时刻(存储为 UTC),例如TIME_ZONE = 'Asia/Shanghai'下30 14 * * *表示每天 14:30 本地时间触发。
内置执行器
| 执行器 | 说明 |
|---|---|
PingPongExecutor |
健康检查,返回 'pong' |
BashScriptExecutor |
执行 Shell 脚本 |
PythonScriptExecutor |
执行 Python 代码 |
SimpleRequestExecutor |
发起 HTTP 请求 |
StatusCheckExecutor |
检测卡住的执行并标记超时(每 5 分钟) |
RetryTimeoutExecutor |
重试超时的执行(每 10 分钟) |
ArchiveExecutor |
归档已完成执行并生成统计(每天凌晨 2 点) |
开发与测试
项目提供 start-test-services.sh 一键启动测试所需的 Redis 与 RabbitMQ 服务。脚本优先检测本地已拉取的镜像(避免依赖 docker.io 上游镜像),也可通过 DJANGO_SIMPLETASK_REDIS_IMAGE / DJANGO_SIMPLETASK_RABBITMQ_IMAGE 环境变量指定,默认使用 6379/5672 端口(被占用时自动选取空闲端口):
./start-test-services.sh start # 启动服务并等待就绪
./start-test-services.sh status # 查看容器与端口状态
./start-test-services.sh stop # 停止并清理容器
运行测试:
python3 -m pytest django_simpletask5_example/tests
测试套件覆盖 Redis、DB、MQ 连接中断后的自动恢复场景,详见 django_simpletask5_example/tests/test_resilience_recovery.py。
测试输出约定:单元测试中对预期会触发的错误(如无效 UUID 校验、毒消息丢弃、消费者断连、执行器加载失败等)所输出的错误日志,必须在测试内屏蔽,避免污染测试输出。统一使用 django_simpletask5_example.tests.suppress_expected_errors 上下文管理器(整体禁用日志)或 patch.object(..., 'logger')(断言同时校验日志)进行屏蔽。
架构
Task 模型变更 → Django 信号 → 创建 TaskExecution 并发布到消息队列
↓
Worker 消费消息 → 获取分布式锁 → 加载执行器 → 执行并保存结果
↓
失败时自动重试/拒绝时重新入队,完成后归档
Releases
0.2.7
- 运行中任务回收机制重构 — 取消原先按
expire_time将 running 直接判为超时的语义;StatusCheckExecutor仅兜底回收长期未被 Worker 获取的 pending。新增core/reclaim,由 crontab 每 ~5 秒基于执行锁探测运行中任务:锁已过期(执行器异常退出)即快速重投或直接置失败,降低孤儿任务滞留时间 - 重试不再搁浅任务 — 失败重试先落库
retry再发布重试消息,发布失败则 requeue 当前消息,杜绝任务永久停在无人处理的retry状态;KombuQueue 依据 AMQPx-death上限丢弃毒消息,避免无限 requeue 造成的队头阻塞与空转 - 消费者重连提示 — Worker 断连重连成功后输出
Consumer connection restored恢复日志,便于运维确认消费通道已恢复 - 安全与健壮性加固 — broker 账号/密码按百分号编码、脚本白名单按路径段边界匹配、signals 快照按实例隔离,覆盖相关回归测试
- 新增测试约定 — 单元测试中对预期触发(如无效 UUID 校验、毒消息丢弃、消费者断连、执行器加载失败等)产生的错误日志必须屏蔽输出,统一使用
django_simpletask5_example.tests.suppress_expected_errors或替换 logger,保证测试输出干净可读;本次同步清理全部此类预期错误日志
0.2.6
- 新增
canceled状态 — 取消已发布但未执行的任务改为独立的终态canceled(原复用failed),Worker 消费到已取消消息时直接跳过;Admin 取消操作、归档/清理/统计(新增canceled_count)、Dashboard 状态图均纳入该状态,取消与执行失败可清晰区分 - WorkerSupervisor 守护线程 — 统一的 Worker 线程保活机制,周期性心跳并自动重启异常退出的 Worker 线程,保持目标 Worker 数量,提升长时间运行稳定性
- 韧性/恢复测试套件 — 新增
test_resilience_recovery.py,模拟 Redis、DB、MQ 连接中断后的自动恢复:Worker 启动时 Redis 注册失败重试、消费者断连自动重连、锁引擎中断后消息重入队并在 Redis 恢复后继续处理、DB 瞬断重入队、执行器异常标记 failed/retry、crontab 调度器 Redis 恢复检测等 - 测试套件稳定性修复 —
test_worker_registry隔离全局_active_workers消除跨用例污染;test_integration动态探测本地可用的 Redis/RabbitMQ 镜像(支持DJANGO_SIMPLETASK_REDIS_IMAGE/DJANGO_SIMPLETASK_RABBITMQ_IMAGE环境变量覆盖),并改为真实 AMQP 握手探测 RabbitMQ 就绪,避免新建容器连接被重置的偶发失败 - 测试服务脚本 —
start-test-services.sh默认使用 6379/5672 端口并动态检测本地镜像,提供start/stop/status子命令,便于一键启动测试依赖
0.2.5
- 修复 0010 迁移双主键冲突 — 修正迁移
0010_taskexecution_id_alter_taskexecution_execution_id的操作顺序:先移除execution_id的主键约束(改为unique=True),再添加自增id主键,避免在已有数据库上升级时因同时存在两个主键而失败。已有 0.2.4 数据可正常升级,id自动回填,execution_id保持唯一
0.2.4
- 破坏性变更:
TaskExecution主键由execution_id(UUID)改为默认自增id,execution_id变为unique=True的唯一标识字段。已有数据库需要执行迁移0010_taskexecution_id_alter_taskexecution_execution_id - 定期存活日志 — Worker 与 Crontab 主循环在长时间无消息/无任务时按
DJANGO_SIMPLETASK_LOG_INTERVAL(默认 300 秒)定期输出存活日志,并统一带时间戳格式,便于运维确认进程未死循环 - MQ 日志降噪 — MQ 消费/发布异常及 Redis 连接失败改为
error/warning级别且无堆栈(默认按间隔节流),避免刷屏;提供--debug时输出完整堆栈。注:此处所指的"锁续期失败"是 0.2.4 版本引入的自定义锁续期机制,自 globallock>=0.2.0 起该机制已移除,改由 globallock 内置看门狗在 lockman/lock 作用域内自动续期,无需(也无从)单独记录续期失败日志 - 新增配置 —
DJANGO_SIMPLETASK_LOG_INTERVAL控制死循环存活日志与异常节流日志的输出间隔(秒)
0.2.3
- Cron 时区修复 — cron 表达式改为按 Django
TIME_ZONE配置的本地时间解释,不再固定按 UTC 计算。USE_TZ=True时使用timezone.localtime(now())作为 croniter 基准时间,非 UTC 时区下定时任务不再偏移;USE_TZ=False保持原有本地时间行为不变
0.2.2
- 抑制异步消息 — 新增
_skip_publish机制,创建/更新/删除 Task 时可通过skip_publish()方法、save(skip_publish=True)、构造时传入skip_publish=True等方式禁止触发 RabbitMQ 消息通知 - 批量创建触发异步消息 —
TaskQuerySet+TaskManager接管bulk_create,批量创建 Task 后自动生成 TaskExecution 并推送 MQ;子类覆盖objects时发出警告提示 - 分布式锁默认引擎 — 默认使用
DjangoRedisGlobalLock,复用 DjangoCACHES['default']配置(需为 Redis 缓存),无需手动配置DJANGO_SIMPLETASK_LOCK_CONFIG - Redis key 命名统一 — 所有 Redis key 格式规范为
{项目前缀}:django-simpletask5:{用途}:{id},项目前缀从ROOT_URLCONF自动提取,也可通过DJANGO_SIMPLETASK_KEY_PREFIX显式覆盖 - CronJob 管理增强 — 新增
app_label字段(自动检测),列表页展示并支持过滤 - 内置 CronJob 重命名 — Status Check → Task Timeout Detection,Retry Timeout → Retry Failed Tasks,Ping/Pong → Health Check 等
- Dashboard 重做 — 新增今日失败卡片,成功率按 90%/80% 分绿色/橙色/红色三级,卡片颜色统一使用 info 蓝色
- Pending/Running 显示修复 — 改为展示
pending / running双值而非仅 pending - 成功率显示优化 —
(22/25)与88%同行显示,字号缩小颜色淡化,高度对齐 - Redis 不可达降级 — Worker 数量显示
-而非0,区分无数据与异常 - 按钮优化 — 改用 Font Awesome 图标,统一
min-width对齐,title 改为中文描述 - 翻译补充 — 补充 Run/Activate/Deactivate 及所有内置 CronJob 名称的中文翻译
- 新增依赖 —
django-static-fontawesome用于按钮图标 - 示例应用 — 示例 cronjobs 模块在 AppConfig.ready() 中导入确保注册生效
0.2.1
- 拒绝回退机制 — 新增
RejectException,执行器可抛出该异常拒绝当前任务,消息自动重新入队等待下次调度,不消耗重试次数 - 区分两种超时 — 引入
pickup_timeout_seconds(等待 Worker 消费超时)与timeout_seconds(执行超时)两种超时配置,expire_time改用 Pickup 超时计算,超时检测更精准 - 文档更新 — README 补充拒绝回退机制的使用说明和两种超时的详细解释
0.2.0
- 破坏性变更: 移除
Task模型上的自定义主键,task_id不再是主键字段,改为unique=True的唯一标识字段。所有Task子模型将自动获得 Django 默认的自增id主键。已有数据库需要迁移处理。 - 统计维度调整:
TaskExecutionStat的统计维度由task_model改为executor_class,统一覆盖有模型任务和 Cron 定时任务两种场景。对应 API_update_stats()按executor_class分组聚合。 - Admin 优化: CronJob 列表页移除
executor_class列避免表格撑开;TaskExecutionStat 列表页新增executor_class列、created_at设为只读修复详情页错误。 - 测试修复: 修复
TransactionTestCase+on_commit在 SQLite 共享内存连接下因TestCase残留原子块导致回调不执行的兼容性问题;Redis MQ / RabbitMQ E2E 集成测试全部通过。
0.1.5
- 安全加固 — 新增
DJANGO_SIMPLETASK_ENABLE_SCRIPT_EXECUTORS配置项(默认关闭),启用后才可执行 Python/Shell 脚本;新增can_use_script_executors权限,仅有此权限或超级管理员的用户才能在 Admin 中创建脚本执行器类型的 CronJob - Bug 修复 — 修复
RetryTimeoutExecutor中 lambda 晚绑定导致所有回调引用最后一个记录的问题 - Admin 安全 — 操作按钮(立即执行、启用/禁用)从 GET 参数改为 POST 请求,附带 CSRF 保护
- 性能优化 —
RetryTimeoutExecutor改为批量 update + 按 ID 分发消息推送 - Admin 增强 —
TaskExecutionAdmin新增executor_class过滤、error_message搜索、批量重试失败任务、批量取消待处理任务 - 数据索引 — 添加
(status, expire_time)和(is_active, next_run_time)复合索引;trigger_event、executor_class、task_id、created_at添加单字段索引 - 代码规范 —
bash_script.py、python_script.py、simple_request.py补充缺失的 logger;移除废弃的allow_tags属性
0.1.4
- 修复打包 —
pyproject.toml添加package-data配置,打包时包含 locale po/mo 文件,修复 i18n 翻译不生效的问题 - 补充依赖 —
requirements.txt和pyproject.toml添加缺失的django-checkbox-normalize、django-tabbed-changeform-admin依赖 - 完善文档 — README 新增完整配置说明章节,涵盖消息队列(MQ)、分布式锁、队列路由、归档、Cron、安全等全部配置项
0.1.3
- Admin 界面大升级 — 集成 Tabbed ChangeForm 分页签展示字段,CronJob 后台支持直接执行/启用/禁用操作;TaskExecution 列表优化,execution_id 显示前8字符并支持点击复制,重试次数合并显示,全局字段中文翻译
- 归档管理增强 — 新增
django_simpletask_archive管理命令,支持手动触发归档与过期清理;PurgeExpiredExecutor 定期自动清理过期记录;保留天数可配置;采用流式 AES-GCM 加密,大文件归档更高效 - Dashboard 修复 — 卡片链接使用
reverse正确跳转,无执行数据时成功率显示- - 超时检测优化 — 模型新增
expire_time字段,超时判定更准确高效
0.1.2
- 修复: 修复 Transaction-Message 排序竞态条件,使用
transaction.on_commit确保消息在事务提交后才发送到队列
0.1.1
- 修复: 修复
DjangoSimpletask5Config缺少ready()方法导致信号监听未注册的问题
0.1.0
这是 django-simpletask5 的首个正式版本。核心功能包括:
- 声明式任务模型 — 继承
Task模型即可定义任务,自动处理创建/更新/删除事件 - 信号驱动自动发布 — Django 信号自动拦截模型变更,发布
TaskExecution到消息队列 - Worker 异步执行 — 多 Worker 进程消费消息队列,支持队列路由和优先级
- Cron 定时调度 — 内置 crontab 守护进程,支持代码注册与数据库覆盖
- 重试与超时 — 失败自动重试(指数退避),支持超时检测
- 加密字段 — 敏感数据自动加密存储
- 归档统计 — 已完成执行记录自动归档为加密 JSONL,并生成日统计
- Worker 注册中心 — 基于 Redis 的 Worker 心跳与状态追踪
- Django Admin 集成 — 完整的后台管理界面与仪表盘
许可证
MIT
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file django_simpletask5-0.2.7.tar.gz.
File metadata
- Download URL: django_simpletask5-0.2.7.tar.gz
- Upload date:
- Size: 96.3 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via:
twine/6.2.0 CPython/3.11.9
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
368ac5acc9d09ee2c2d0f9a5b3e95615496d21cdfa0c6f1f14873828fd6662d1
|
|
| MD5 |
4b720b6ffbd4b3e53a3e75e5143a71dd
|
|
| BLAKE2b-256 |
0b5563071460ac27a707e4f6bb0f69e9e3abc2c9c5b27fbfb21296e61d09a5a6
|
File details
Details for the file django_simpletask5-0.2.7-py3-none-any.whl.
File metadata
- Download URL: django_simpletask5-0.2.7-py3-none-any.whl
- Upload date:
- Size: 83.8 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via:
twine/6.2.0 CPython/3.11.9
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
3df6e346a3358a6c94503d0c990dcff1797a5ce2e0ba72b94db9563ae2dbad9b
|
|
| MD5 |
2c5264dca78bcc9bce592bd19600c1d2
|
|
| BLAKE2b-256 |
9c7f8573ec033356174ba5c2c177953d9a753a6ccd642a00821abb9adb293e9d
|