- 任务调度
- 后端
- 消息队列
【免费下载链接】celery
Distributed Task Queue (development branch)
导读
本文以 Celery 仓库中的 docs/reference/celery.beat.rst API 参考文档为骨架,深入剖析其核心对象celery.beat模块(源码位于 celery/beat.py)。该模块是celery beat周期任务调度器的实现主体,负责按beat_schedule配置在固定间隔内将任务发送到 Broker。读完本文,你将掌握ScheduleEntry、Scheduler、PersistentScheduler、Service、EmbeddedService等核心类的职责与协作方式,理解调度堆(heap)与 tick 循环的工作原理、持久化与同步机制,以及相关配置项、CLI 命令与信号的完整用法。
一、模块定位:celery.beat 是什么
celery.beat是 Celery 的周期任务调度器(periodic task scheduler)模块。正如其模块 docstring 所言——"The periodic task scheduler"——它独立于 worker 运行,按预定义的时间表周期性地把任务发送到消息队列,由 worker 消费执行。
在 docs/userguide/periodic-tasks.rst 中,:program:celery beat`` 被描述为"a scheduler; It kicks off tasks at regular intervals"。默认情况下,调度条目(entries)取自beat_schedule配置项。而 celery/apps/beat.py 则明确写道,它是celery.beat模块的"program-version"——即负责把celery.beat真正跑成一个应用程序,包括安装信号处理器、创建 pidfile、设置进程标题、配置日志等。
celery.beat模块导出的公开 API(见 celery/beat.py 的__all__)为:
SchedulingError:调度过程中抛出的异常ScheduleEntry:调度表中的一个条目Scheduler:周期任务调度器基类PersistentScheduler:基于shelve数据库持久化的调度器Service:Celery 周期任务服务EmbeddedService:嵌入式时钟服务工厂函数
另有未列入__all__的辅助类BeatLazyFunc(惰性函数包装)与命名元组event_t = namedtuple('event_t', ('time', 'priority', 'entry')),后者是调度堆中每个元素的载体。
二、ScheduleEntry:调度表的基本单元
ScheduleEntry(源码见 celery/beat.py)代表调度表中的一条记录,即"某个任务按照某个时间表周期性执行"的完整描述。
2.1 字段定义
| 字段 | 类型 | 说明 |
|---|---|---|
name | str | 调度条目的名称(唯一标识) |
task | str | 要调度的任务名 |
schedule | ~celery.schedules.schedule | 时间表对象(如crontab、timedelta等,经maybe_schedule统一转换) |
args | Tuple | 传给任务的位置参数 |
kwargs | Dict | 传给任务的关键字参数 |
options | Dict | 任务执行选项 |
last_run_at | datetime | 上次计划执行的时间 |
total_run_count | int | 该条目已被调度的总次数 |
relative | bool | 时间是否相对于服务器启动时刻 |
构造函数签名(celery/beat.py):
def __init__(self, name=None, task=None, last_run_at=None, total_run_count=None, schedule=None, args=(), kwargs=None, options=None, relative=False, app=None):注意几个实现细节:
schedule通过maybe_schedule(schedule, relative, app=self.app)转换,因此既可以直接传~celery.schedules.schedule实例,也可以传整数秒数或datetime.timedelta(详见 celery/schedules.py 的schedule类)。last_run_at缺省时取self.schedule.now()(若 schedule 存在)否则取self.app.now()。- 该类的
__iter__返回vars(self).items(),支持dict(self)形式的序列化重建(_next_instance正是这样构造下一个实例的)。
2.2 核心行为
_next_instance()/next()/__next__:返回一个新实例,仅更新last_run_at(默认取当前时间)并将total_run_count加 1。调度器每次真正派发一个条目后,就用它替换旧条目,从而推进调度状态。is_due():委托给self.schedule.is_due(self.last_run_at),返回(is_due, next_time_to_run)二元组——是否到期,以及距离下次运行的时间(秒)。update(other):只更新"可编辑字段"task、schedule、args、kwargs、options,用于磁盘上的旧条目与配置中新增定义合并。__eq__/editable_fields_equal:相等性只比较上述可编辑字段;而__lt__则退回id(self) < id(other),正如注释所解释的——在调度堆中,排序由元组(time, priority, entry)的前两个成员决定,堆顶竞争最后才轮到比较 entry,因此此时顺序"随机即可"。
从源码结构看,ScheduleEntry的"可编辑字段"与"运行状态字段(last_run_at / total_run_count)"的分离是刻意设计:前者用于合并外部配置更新,后者用于跟踪调度进度。对应的单元测试位于 t/unit/app/test_beat.py,其中test_next、test_is_due、test_update、test_repr、test_reduce、test_lt等用例直接覆盖了上述行为。
三、BeatLazyFunc:调度参数中的惰性求值
BeatLazyFunc(celery/beat.py)用于在beat_schedule中声明"发送任务前才调用"的惰性函数。其典型场景是:某些参数(如当前时间)必须在任务真正派发那一刻才求值,而不能在配置加载时固定下来。
官方 docstring 给出的示例:
beat_schedule = { 'test-every-5-minutes': { 'task': 'test', 'schedule': 300, 'kwargs': { "current": BeatCallBack(datetime.datetime.now) } } }实现要点:
- 构造时接收
func及预置的args/kwargs,保存在self._func_params。 - 定义
__call__与delay(),二者都会在调用时执行self._func(*args, **kwargs)返回真实值。
它真正被消费的位置在Scheduler.apply_async内部(见下文 4.4 节):_evaluate_entry_args与_evaluate_entry_kwargs(celery/beat.py)会遍历条目的args/kwargs,把其中所有BeatLazyFunc实例替换为其调用结果。t/unit/app/test_beat.py的test_beat_lazy_func对该机制进行了验证。
四、Scheduler:调度器核心
Scheduler(celery/beat.py)是所有调度器的基类,实现了完整的调度主循环、到期判定、任务派发与同步机制。celery beat程序可能出于内省目的多次实例化本类,此时会传入lazy=True——因此文档特别强调子类在lazy参数下必须保持幂等。
4.1 关键类属性
| 属性 | 默认值 | 说明 |
|---|---|---|
Entry | ScheduleEntry | 使用的条目类型,子类可覆写 |
max_interval | DEFAULT_MAX_INTERVAL = 300(5 分钟) | 两次检查调度表之间的最大睡眠秒数 |
sync_every | 3 * 60(3 分钟) | 多久强制同步一次调度表 |
sync_every_tasks | None | 每派发多少任务强制同步一次(默认不启用) |
4.2 初始化与配置解析
构造函数(celery/beat.py):
def __init__(self, app, schedule=None, max_interval=None, Producer=None, lazy=False, sync_every_tasks=None, **kwargs): self.app = app self.data = maybe_evaluate({} if schedule is None else schedule) self.max_interval = (max_interval or app.conf.beat_max_loop_interval or self.max_interval) self.Producer = Producer or app.amqp.Producer ... self.sync_every_tasks = ( app.conf.beat_sync_every if sync_every_tasks is None else sync_every_tasks) if not lazy: self.setup_schedule()注意配置优先级的实际体现:max_interval依次取构造函数参数 →app.conf.beat_max_loop_interval→ 类属性默认值 300 秒;sync_every_tasks缺省时取app.conf.beat_sync_every。这直接对应 celery/app/defaults.py 中beat命名空间的配置定义。
setup_schedule()(celery/beat.py)做两件事:
install_default_entries(self.data):若result_expires已配置且后端不支持自动过期,则自动注册celery.backend_cleanup任务(默认crontab('0', '4', '*'),即每天凌晨 4 点,options带expires: 12 * 3600),用于清理过期结果。merge_inplace(self.app.conf.beat_schedule):将beat_schedule配置合并进调度表。
merge_inplace(celery/beat.py)的合并语义值得注意:
- 计算配置键集合
B与现有调度键集合A的对称差A ^ B,从调度表中移除"磁盘上存在但配置里已删除"的条目; - 对
B中的每个键:若已存在则调用schedule[key].update(entry)仅刷新可编辑字段;否则新建条目。这样既能应用配置变更,又不丢失last_run_at/total_run_count等运行状态。
4.3 调度堆(heap)与 tick 循环
Scheduler的核心数据结构是一个最小堆self._heap,堆元素为event_t(time, priority, entry)三元组,按触发时间戳排序。
populate_heap()(celery/beat.py):遍历self.schedule.values(),对每个条目计算is_due, next_call_delay = entry.is_due(),若已到期则时间戳按 0 计算,否则按next_call_delay计算,随后heapify。时间戳由_when()生成:将last_run_at归一化为 UTC 后加上微秒与经过adjust()(默认减去 0.010 秒漂移补偿)的延迟值。
tick()(celery/beat.py)——每次调用执行"一个到期任务",是调度的主循环体,返回"下次调用前应等待的秒数"建议。其流程如下:
- 若堆为空或调度表发生变化(
schedules_equal比较键集合与各条目可编辑字段),则重建堆。 - 取堆顶事件:若
event[0] > now,说明最早的到期时间在未来,返回min(event[0] - now, max_interval)作为睡眠建议。 - 若堆顶已到期:
heappop弹出并校验身份后,reserve(entry)(即self.schedule[entry.name] = next(entry),推进运行状态)并apply_entry派发任务,然后把next_entry按新的触发时间重新入堆,返回 0(立即继续下一轮)。 - 若堆顶已到期但条目自身要求"稍后重试":将其按重试时间重新入堆(
reschedule_delay),避免该条目长期占据堆顶、饿死后续条目——这一点在代码注释中明确引用了https://github.com/celery/celery/issues/7649的问题背景。
tick的这些分支在 t/unit/app/test_beat.py 中有大量测试覆盖,例如test_due_tick、test_due_tick_returns_delay_when_heap_top_changed、test_pending_tick、test_honors_max_interval、test_not_due_top_entry_is_rescheduled_behind_due_entry,以及针对错过 cron 截止期限的test_tick_dispatches_missed_cron_within_deadline_non_uniform等。
4.4 任务派发:apply_entry 与 apply_async
apply_entry(celery/beat.py)负责实际"发送到期任务",记录 info 级日志Scheduler: Sending due task %s (%s),派发失败时记录 error 日志并继续(不会让整个 beat 崩溃)。
真正的工作在apply_async(celery/beat.py)中完成,要点如下:
- 先
reserve推进时间戳与计数(注释强调要在真正执行前完成,避免异常导致"永远重复调度")。 - 对
entry.args/entry.kwargs执行惰性求值(见第三节BeatLazyFunc)。 - 给消息注入自定义 header:
options.setdefault('headers', {})['celery_beat_task'] = True,用于标识消息来源于 Celery Beat。 - 若任务已注册,走
task.apply_async(...);否则退回self.send_task(...)(即app.send_task)。 - 派发后递增
_tasks_since_sync,若should_sync()为真则执行_do_sync()同步调度表。 - 任何异常都会被包装为
SchedulingError重新抛出,消息形如Couldn't apply scheduled task {name}: {exc}。
apply_async的行为在测试中有直接印证:test_apply_async_uses_registered_task_instances、test_apply_async_with_null_args、test_apply_async_with_null_args_set_to_none等用例(t/unit/app/test_beat.py)。
4.5 同步机制:should_sync / sync / close
should_sync()(celery/beat.py):满足任一条件即需同步——距上次同步超过sync_every秒,或(启用了sync_every_tasks时)自上次同步以来派发任务数达到阈值。_do_sync():调用self.sync()并重置_last_sync与_tasks_since_sync。- 基类的
sync()为空实现(内存型调度器无需落盘),close()也仅调用sync();PersistentScheduler会覆写二者。
4.6 连接与 Producer
connection与producer均为cached_property(celery/beat.py):
@cached_property def connection(self): return self.app.connection_for_write() @cached_property def producer(self): return self.Producer(self._ensure_connected(), auto_declare=False)_ensure_connected()通过connection.ensure_connection(_error_handler, self.app.conf.broker_connection_max_retries)建立连接,连接失败时以beat: Connection error: %s. Trying again in %s seconds...记录错误并重试,重试上限由broker_connection_max_retries控制。
五、PersistentScheduler:基于 shelve 的持久化调度器
PersistentScheduler(celery/beat.py)是默认使用的调度器(配置beat_scheduler的默认值即为'celery.beat:PersistentScheduler',见 celery/app/defaults.py),它把调度表持久化到shelve数据库中,默认文件名celerybeat-schedule。
5.1 文件与损坏恢复
persistence = shelve,known_suffixes = ('', '.db', '.dat', '.bak', '.dir')——_remove_db()会按这些后缀逐一尝试删除(用platforms.ignore_errno(errno.ENOENT)忽略文件不存在)。_open_schedule()以writeback=True打开数据库。setup_schedule()打开数据库并读取keys()以触发潜在损坏错误;若失败,则记录Removing corrupted schedule file %r: %r并删除数据库重建——代码注释特别提到 bsddb 的DBPageNotFoundError这类"打开成功但首次取键才报错"的损坏场景。
5.2 时区 / UTC 变更重置
持久化数据库中会额外保存__version__、tz、utc_enabled三个元数据字段:
- 若存储的
tz与当前app.conf.timezone不同,或存储的utc_enabled与app.conf.enable_utc不同,则警告Reset: Timezone changed from %r to %r/Reset: UTC changed from %s to %s并clear()整个数据库——因为时间表语义依赖时区,变更时必须重置以免按错误时区触发。 _create_schedule()中的升级逻辑对旧版本数据库做了兼容:缺少__version__字段(2.2.2 之前)→ 重置;缺少tz(3.0.8 之前)→ 重置;缺少utc_enabled(3.0.9 之前)→ 重置。
5.3 覆盖的接口
schedule属性改为self._store['entries']的读写。sync()调用self._store.sync();close()在sync()后关闭self._store。info返回f' . db -> {self.schedule_filename}',用于启动横幅展示。
六、Service 与 EmbeddedService:把调度器跑起来
6.1 Service
Service(celery/beat.py)把调度器封装为可启动/停止的服务,是celery beat命令行与嵌入模式共用的运行时外壳:
scheduler_cls = PersistentScheduler(默认),构造函数中的schedule_filename缺省取app.conf.beat_schedule_filename,max_interval缺省取app.conf.beat_max_loop_interval。- 使用两个
threading.Event:_is_shutdown与_is_stopped协调启停。 start(embedded_process=False)的主循环:
while not self._is_shutdown.is_set(): interval = self.scheduler.tick() if interval and interval > 0.0: debug('beat: Waking up %s.', humanize_seconds(interval, prefix='in ')) time.sleep(interval) if self.scheduler.should_sync(): self.scheduler._do_sync()启动时发送beat_init信号;若为嵌入进程还发送beat_embedded_init信号并将进程标题设为celery beat;捕获KeyboardInterrupt/SystemExit后置_is_shutdown,最终self.sync()关闭调度器并置_is_stopped。
stop(wait=False)置_is_shutdown,可选阻塞等待_is_stopped。get_scheduler(lazy=False, extension_namespace='celery.beat_schedulers'):通过load_extension_class_names加载celery.beat_schedulers命名空间下的第三方调度器别名,再用symbol_by_name解析scheduler_cls——这是自定义调度器(如数据库调度器)被beat_scheduler配置引用的解析入口。
6.2 EmbeddedService
EmbeddedService(celery/beat.py)是返回嵌入式时钟服务的工厂函数:
def EmbeddedService(app, max_interval=None, **kwargs): if kwargs.pop('thread', False) or _Process is None: # Need short max interval to be able to stop thread # in reasonable time. return _Threaded(app, max_interval=1, **kwargs) return _Process(app, max_interval=max_interval, **kwargs)- 默认使用多进程
_Process(基于billiard的Process),run()中会reset_signals、关闭标准输入输出与日志文件描述符,并调用service.start(embedded_process=True)。 - 传
thread=True(或平台不支持多进程)时回退到_Threaded,其max_interval被强制设为 1 秒以便线程能在合理时间内被停止;线程名为Beat,且为 daemon 线程。
这解释了嵌入式模式(如 Django 中通过app.Beat或直接调用EmbeddedService在应用进程内运行 beat)的两种形态,也对应 celery/signals.py 中beat_init与beat_embedded_init两个信号的存在意义。
七、配置项速查(beat 命名空间)
以下配置定义于 celery/app/defaults.py 的beat命名空间:
| 配置键 | 默认值 | 类型 | 作用 |
|---|---|---|---|
beat_max_loop_interval | 0 | float | 调度循环两次检查之间的最大睡眠秒数;0 表示使用调度器类默认值(300 秒) |
beat_schedule | {} | dict | 周期任务调度表,键为条目名,值为{task, schedule, args, kwargs, options}字典 |
beat_scheduler | 'celery.beat:PersistentScheduler' | str | 调度器类,可为点路径或注册在celery.beat_schedulers命名空间的别名 |
beat_schedule_filename | 'celerybeat-schedule' | str | 调度数据库文件名 |
beat_sync_every | 0 | int | 每派发 N 个任务强制同步调度表;0 表示关闭(仅按时间同步) |
beat_cron_starting_deadline | None | int | 错过 cron 触发后的补偿截止期限(秒),相关分支见tick()与测试用例test_tick_dispatches_missed_cron_within_deadline_non_uniform |
另外,Scheduler中sync_every = 3 * 60的"按时间同步"默认间隔是类属性而非配置项;beat_schedule的完整写法可参考 docs/userguide/periodic-tasks.rst 中app.conf.beat_schedule的示例(支持crontab与整数秒/timedelta两种 schedule 类型)。
八、命令行:celery beat
8.1 启动命令
celery -A proj beat默认从当前目录读取celerybeat-schedule数据库文件(若不存在则创建)。celery.beat模块在此处的角色是"被驱动方"——CLI 层在 celery/bin/beat.py 中定义,应用层封装在 celery/apps/beat.py 的Beat类,最终调用beat.Service即celery.beat的Service。
8.2 常用选项(Beat Options,见 celery/bin/beat.py)
| 选项 | 说明 |
|---|---|
--detach | 后台守护进程方式运行 |
-s, --schedule | 调度数据库路径,默认取beat_schedule_filename配置(默认celerybeat-schedule),扩展名.db会被追加 |
-S, --scheduler | 使用的调度器类,默认取beat_scheduler配置(即celery.beat:PersistentScheduler) |
--max-interval | 调度迭代之间最大睡眠秒数(int),对应max_interval |
-l, --loglevel | 日志级别,默认WARNING |
-f, --logfile、--pidfile、--uid、--gid、--umask、--workdir | 继承自守护命令的通用选项 |
此外还支持-C(无彩色输出)与额外配置参数透传:ctx.args中未被识别的内容会交给app.config_from_cmdline(ctx.args)解析,解析失败时抛出click.UsageError。
8.3 启动流程与信号处理
Beat.run()(celery/apps/beat.py)依次执行:
- 打印
celery beat v{VERSION_BANNER} is starting.横幅; init_loader():app.loader.init_worker()导入任务模块并app.finalize();- 设置进程标题
celery beat; start_scheduler():若指定pidfile则platforms.create_pidlock;构造Service并打印启动信息(STARTUP_INFO_FMT展示 broker、loader、scheduler、logfile、maxinterval 等);设置日志;可选设置全局 socket 超时(默认 30 秒);安装同步信号处理器后service.start()。
install_sync_handler(celery/apps/beat.py)为SIGTERM与SIGINT安装处理器:收到信号时先service.sync()保存调度表,再抛SystemExit——确保关闭时调度状态(如last_run_at、total_run_count)不丢失。
九、信号钩子
celery.beat相关信号定义于 celery/signals.py:
beat_init:在Service.start()中、主循环开始前发送(signals.beat_init.send(sender=self)),可用于初始化资源。beat_embedded_init:仅在嵌入式进程模式(embedded_process=True)下发送,用于区分独立进程与嵌入应用进程两种运行形态。
十、扩展调度器与测试验证
10.1 自定义调度器
Service.get_scheduler的extension_namespace='celery.beat_schedulers'表明,第三方调度器(例如将调度表存入数据库的django_celery_beat.schedulers:DatabaseScheduler)通过注册到该命名空间,即可用配置或-S选项引用。自定义调度器通常继承Scheduler,覆写setup_schedule、sync、close与schedule属性即可(参考PersistentScheduler的模式)。
10.2 测试支撑
模块级行为在 t/unit/app/test_beat.py 中有完整覆盖,可作为理解实现语义的"活文档",代表性用例包括:
- 惰性参数求值:
test_beat_lazy_func; - 条目推进:
test_next、test_reduce; - tick 分支:
test_due_tick、test_pending_tick、test_honors_max_interval、test_ticks_schedule_change、test_not_due_top_entry_is_rescheduled_behind_due_entry; - 同步计数:
test_should_sync、test_sync_task_counter_resets_on_do_sync; - 默认条目与合并:
test_install_default_entries、test_merge_inplace; - 错过 cron 的截止期限补偿:
test_tick_dispatches_missed_cron_within_deadline_non_uniform等。
结语
celery.beat模块以"条目(ScheduleEntry)+ 调度器(Scheduler)+ 服务(Service)"三层结构,将周期任务调度抽象为清晰可扩展的组件:堆驱动的 tick 循环保证调度效率,持久化与同步机制保证状态不丢失,beat_schedule配置与celery beat命令提供开箱即用的使用入口,而beat_init/beat_embedded_init信号与celery.beat_schedulers扩展命名空间则为其在框架内的定制留足空间。理解这一模块,是深入使用与扩展 Celery 周期任务能力的关键一步。
- 任务调度
- 后端
- 消息队列
【免费下载链接】celery
Distributed Task Queue (development branch)
相关推荐
Celery 周期性任务(Periodic Tasks)完整指南:beat 调度器、Crontab 与 Solar 调度实战
Celery 周期性任务(Periodic Tasks)完整指南:beat 调度器、Crontab 与 Solar 调度实战 导读 本文以 Celery 官方用
任务调度后端消息队列雀魂数据分析终极指南:如何用免费工具快速提升麻将水平
雀魂数据分析终极指南:如何用免费工具快速提升麻将水平 想要从雀魂麻将新手成长为高手吗?雀魂牌谱屋(amae koromo)就是你需要的秘密武器!这款完全免费的开
CANNAscend人工智能任务调度celery beat 命令完全指南:配置、调度器与源码级原理剖析
celery beat 命令完全指南:配置、调度器与源码级原理剖析 导读 : celery beat 是 Celery 分布式任务队列内置的周期任务调度器,负责
任务调度后端消息队列
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考