Skip to main content

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 默认配置项(DATABASESMIDDLEWARETEMPLATESROOT_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 仅用于开发调试;生产必须改为 rabbitmqredis,否则消息不持久化、跨进程 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 默认自增 idTask 不再内置 task_id 字段,业务标识由业务模型自行定义(自定义主键或自然键均可)
  • 保存 / 删除 / 批量创建自动触发事件,详见「事件处理机制」
  • 事件上下文、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,执行完成变为 successresult 保存执行器返回值
  • Django Admin 的 Dashboard 可查看执行统计与 Worker 状态

事件处理机制

内置事件:create / update / delete

框架通过 Django 信号自动监听所有 Task 子类的保存与删除,无需手动发布:

事件 触发时机 执行器可见的上下文
create 新增实例(save() / bulk_create() task_idtrigger_eventdata(全字段快照)
update 字段值发生变化后保存 task_idtrigger_eventdataold_datanew_datachanged_fields
delete 调用 instance.delete() task_idtrigger_eventdata

执行器通过 execution.get_context_dict() 获取以上上下文。其中 task_id 是业务对象主键的字符串表示(如 "42"、UUID 字符串或自定义键);业务对象反查见「模型与生命周期」。

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'

这样实现删除操作与外部系统清理的解耦,且支持失败重试。

confirm_delete()execution.get_task() 都会通过 TaskExecution.task_id(业务主键字符串)反查业务对象,反查逻辑由业务模型上的 get_task_instance() 提供(见「模型与生命周期」)。

事件 → 执行器映射

每个事件消息最终要由一个执行器处理。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)→ defaultcreate → 自动派生 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 字符串);返回 Noneresult 保持为空

批量任务

重写 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 便于进度展示。

异常与状态流转

执行器抛出的异常由框架统一处理:

情况 结果
无异常 状态 successresult 保存返回值
RejectException 拒绝回退:状态复位 pending,消息重新入队,不消耗重试次数
其他异常且 retry_count < max_retries 状态 retryretry_count + 1,消息重新入队
其他异常且重试次数已用完 状态 failederror_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 状态恢复为 pendingstarted_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

  • 主键为 Django 默认自增 id,无需自定义;不内置 task_id 字段,业务标识由业务模型按需定义
  • 提供 trigger(event, extra_context)skip_publish()get_executor_class_path(event)get_queue_name(event) 等方法
  • 保存、删除、批量创建自动触发事件(见「事件处理机制」)

业务对象反查:get_task_instance()

TaskExecution.task_id 存的是业务对象主键的字符串表示(CharField)。需要反查业务对象时(execution.get_task() / execution.confirm_delete()),框架先通过 task_model 解析出模型类,再调用该类的 get_task_instance(task_id) 完成查询。

默认实现在 Task 基类中按主键 pk 查询,并通过内置的 _coerce_task_pk() 兼容常见主键类型:

  • AutoField / BigAutoField / SmallAutoField — 字符串转为 int
  • UUIDField — 字符串转为 uuid.UUID
  • CharField / TextField 等字符串主键 — 原样使用
task = execution.get_task()          # 返回业务对象(反查不到返回 None)
execution.confirm_delete()          # 反查并真实删除业务对象

业务模型使用自定义主键或非主键自然键时,重载 get_task_instance() 即可,框架会自动改走你的实现:

class OrderTask(Task):
    code = models.CharField(max_length=64, unique=True)   # 非主键自然键

    @classmethod
    def get_task_instance(cls, task_id):
        if task_id is None:
            return None
        return cls.objects.filter(code=task_id).first()

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 记录标记为 timeoutRetryTimeoutExecutortimeout 记录重新入队重试(执行超时消耗重试次数,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.8

  • 移除 Task 冗余的 task_id 字段TaskExecution.task_id 由 UUID 改为 CharField(max_length=255),直接用于存放业务对象主键字符串(str(instance.pk)),不再维护一份与业务主键重复的 UUID。消息上下文中的 task_id 字段名不变,值语义从 UUID 变为业务主键字符串(如 "42"
  • 业务对象反查下沉到模型层 — 新增 Task.get_task_instance(task_id) 类方法作为默认反查实现(按 pk 查询,内置 _coerce_task_pk() 兼容 int / UUID / 字符串主键);TaskExecution.get_task() / confirm_delete() 改为解析 task_model 后委托给该方法。业务模型使用自定义主键或自然键时重载 get_task_instance() 即可
  • 新增迁移django_simpletask5.0013_alter_taskexecution_task_id(字段改 CharField)、django_simpletask5_example.0003_remove_ordertask_task_id
  • 破坏性变更: 依赖 TaskExecution.task_id 为 UUID 类型的外部消费方需适配为业务主键字符串;业务 Task 子类若曾使用内置 task_id 字段需移除

0.2.7

  • 运行中任务回收机制重构 — 取消原先按 expire_time 将 running 直接判为超时的语义;StatusCheckExecutor 仅兜底回收长期未被 Worker 获取的 pending。新增 core/reclaim,由 crontab 每 ~5 秒基于执行锁探测运行中任务:锁已过期(执行器异常退出)即快速重投或直接置失败,降低孤儿任务滞留时间
  • 重试不再搁浅任务 — 失败重试先落库 retry 再发布重试消息,发布失败则 requeue 当前消息,杜绝任务永久停在无人处理的 retry 状态;KombuQueue 依据 AMQP x-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)改为默认自增 idexecution_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,复用 Django CACHES['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_eventexecutor_classtask_idcreated_at 添加单字段索引
  • 代码规范bash_script.pypython_script.pysimple_request.py 补充缺失的 logger;移除废弃的 allow_tags 属性

0.1.4

  • 修复打包pyproject.toml 添加 package-data 配置,打包时包含 locale po/mo 文件,修复 i18n 翻译不生效的问题
  • 补充依赖requirements.txtpyproject.toml 添加缺失的 django-checkbox-normalizedjango-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

django_simpletask5-0.2.8.tar.gz (99.9 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

django_simpletask5-0.2.8-py3-none-any.whl (86.0 kB view details)

Uploaded Python 3

File details

Details for the file django_simpletask5-0.2.8.tar.gz.

File metadata

  • Download URL: django_simpletask5-0.2.8.tar.gz
  • Upload date:
  • Size: 99.9 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.11.9

File hashes

Hashes for django_simpletask5-0.2.8.tar.gz
Algorithm Hash digest
SHA256 80e74f22b10225a7046d682c461d45e3c936b6b30e86cad5f329a5c671298c9c
MD5 b4028ac92cc4ef042dfa7f693af30e90
BLAKE2b-256 db0473cd60abf44e2cd42d3c8c169d3f3bb94201d124898689a24f25f59fe2e3

See more details on using hashes here.

File details

Details for the file django_simpletask5-0.2.8-py3-none-any.whl.

File metadata

File hashes

Hashes for django_simpletask5-0.2.8-py3-none-any.whl
Algorithm Hash digest
SHA256 a43a6d3a646a1291cd615acf274ca52c7b12d855cdf6c4e97c2325720224bcab
MD5 4e0c5792923335f92f4e6a628af23d62
BLAKE2b-256 a67060825458c7fad78e06ae88054e9ac8f39f7bb27f55f11ac337053f58438d

See more details on using hashes here.

Release history Release notifications | RSS feed

This release

0.2.8 This release

2 files

0.2.7

2 files

0.2.6

2 files

0.2.5

2 files

0.2.4

2 files

0.2.3

2 files

0.2.2

2 files

0.2.1

2 files

0.2.0

2 files

0.1.4

2 files

0.1.3

2 files

0.1.2

2 files

0.1.1

2 files

0.1.0

2 files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page