使用sqlalchemy_celery_beat调度器
对超长时间eta和很多数量的待执行任务友好
还有一种方法是使用RabbitMQ延迟插件,待研究
因为celery-sqlalchemy-beat很久没维护依赖已经过时(不支持sqlalchemy2.0版本),所以使用包名:sqlalchemy-celery-beat
先安装
uv pip install gevent celery[redis] sqlalchemy
再安装
uv pip install sqlalchemy-celery-beat
周期任务不执行的可能原因
假设有一个一天执行一次的周期任务,因在调试开发中经常改时间,celery_crontabschedule表中的last_run_at字段会记录当前任务的上次执行时间,total_run_count执行次数,如果这个值未超过24小时(设置时间的24小时内),会被认为这个任务当天已执行,因此它可能会跳过当次执行,因此需要清空last_run_at、total_run_count字段
代码示范(推荐后面的稳健行+重连数据库)
celery_app.py:
# celery任务实列,调度工具
from celery import Celery
from celery.signals import beat_init,worker_shutdown #进程生命周期管理
from sqlalchemy_celery_beat.models import (
PeriodicTask, # 任务表
ClockedSchedule, # 一次性任务
PeriodicTaskChanged, # 任务变更
IntervalSchedule, # 间隔任务
CrontabSchedule, # 周期任务
ModelBase, # 继承模型基类
)
from sqlalchemy import Column, Integer, String, DateTime,inspect
from sqlalchemy_celery_beat.session import SessionManager
import json
import random
from datetime import datetime, timedelta
from contextlib import contextmanager
REDIS_URL = "redis://:foobare122Edsd@192.168.1.11:6379/1"
# 数据库连接添加减少意外关闭概率
POSTGRES_URL = "postgresql://postgres:123456@192.168.1.13:5432/celery?sslmode=require&keepalives=1&keepalives_idle=30&keepalives_interval=5&keepalives_count=3"
# SQLITE_URL = "sqlite:///celery_beat.db" # 使用sqlite数据库
# 创建celery实例
app = Celery('test_app', broker=REDIS_URL,backend=None)
# 配置celery
app.conf.update(
{
'beat_dburi': POSTGRES_URL, # 设置beat数据库uri
'beat_schema': None, # 设置beat数据库schema
'timezone':'Asia/Singapore', # 时区
'enable_utc': False, # 不使用UTC时区
'include': ['tasks'], # 导入多个任务模块,用来自动发现任务,避免循环导入
# 全局任务超时配置,如果周期任务超过了这个时间,则任务会被删除
'task_annotations': {
'*': {
'expires': 1800, # 30分钟后过期(相对时间)
}
},
}
)
# 创建一个 session 管理器
session_manager = SessionManager()
@contextmanager
def get_session():
"""数据库会话上下文管理器"""
session = session_manager.session_factory(POSTGRES_URL)
try:
yield session
session.commit()
except Exception:
session.rollback()
raise
finally:
session.close()
# 定义一个已完成一次性任务表
class ExecutedOneTasks(ModelBase):
__tablename__ = 'celery_executed_one_tasks'
id = Column(Integer, primary_key=True)
task_name = Column(String(255), nullable=False)
task = Column(String(255), nullable=False)
task_args = Column(String(255))
# 执行时间
last_run_at = Column(DateTime)
# # 任务调度类型
# discriminator = Column(String(255), nullable=False)
# 初始化数据库
def init_db():
"""初始化数据库,清理所有周期任务"""
with get_session() as session:
# 创建一个已执行的表
inspector = inspect(session.bind)
if 'celery_executed_one_tasks' not in inspector.get_table_names():
# 表不存在,创建它
session.execute(
ExecutedOneTasks.__table__.create(bind=session.bind)
)
# 确保 PeriodicTaskChanged 记录存在(仅初始化时调用一次)
change_record = session.query(PeriodicTaskChanged).first()
if not change_record:
change_record = PeriodicTaskChanged()
session.add(change_record)
# 删除所有过期超过1天的celery_clockedschedule任务
session.query(ClockedSchedule).filter(
ClockedSchedule.clocked_time < datetime.now() - timedelta(days=1)
).delete()
# 删除所有间隔任务
session.query(PeriodicTask).filter_by(
discriminator='intervalschedule'
).delete()
# 删除所有周期任务
session.query(PeriodicTask).filter_by(
discriminator='crontabschedule'
).delete()
# 删除所有间隔调度
session.query(IntervalSchedule).delete()
# 删除所有周期调度
session.query(CrontabSchedule).delete()
# 生成唯一任务名称
def generate_task_name(task_name):
"""
生成唯一任务名称
参数:
task_name: 任务名称前缀
返回格式: prefix_timestamp_random
"""
timestamp = int(datetime.now().timestamp())
random_suffix = random.randint(100, 999) # 3位随机数
return f"{task_name}_{timestamp}_{random_suffix}"
# 创建一个一次性任务对象
class TaskOneObj:
def __init__(self,task:str,clocked_time:datetime,task_args:list=[]):
self.task = task # 任务函数名(任务名)
self.clocked_time = clocked_time # 任务发送到队列时间
self.task_args = task_args # 任务参数
# 创建一次性任务
def create_one_tasks(tasks_obj:list[TaskOneObj]):
"""创建一次性任务(传入自定义的任务对象数组,使用数组批量对数据库性开销好),每个对象的属性如下:
参数:
task: 任务函数名(任务名)
clocked_time: 任务发送到队列时间
task_args: 任务参数
"""
with get_session() as session:
for task_obj in tasks_obj:
# 创建 一个一次性clocked时间调度
clocked_schedule = ClockedSchedule(clocked_time=task_obj.clocked_time)
session.add(clocked_schedule)
session.flush()
#创建一个一次性任务
one_task = PeriodicTask(
schedule_model=clocked_schedule,
name=generate_task_name(task_obj.task),
task=task_obj.task,
args=json.dumps(task_obj.task_args), # 任务位置参数,按顺序传入参数
#kwargs=json.dumps({'priority': 2, 'delay': 10}), # 任务关键字参数,按参数名传入参数,这里是字典
# 任务生效时间,建议不设置,只设置clocked_time就行
# start_time=datetime.now(tz=timezone.utc) + timedelta(minutes=1),
one_off=True, # 设置为一次性任务
enabled=True, # 启用任务
# expires=task_obj.clocked_time + timedelta(seconds=1800),# 任务过期时间(可选),1800秒后任务过期
)
# 提交一次性任务
session.add(one_task)
# 通知调度器有新任务,更新celery_periodictaskchanged表里的一条时间数据
session.execute(
PeriodicTaskChanged.__table__.update().values(
last_update=datetime.now()
)
)
# 创建或更新间隔任务,间隔任务相同频率的可以复用一个调度
def create_interval_tasks(every:int,period:str,task:str):
"""创建间隔任务,通常间隔任务周期任务同一个任务名只存在一个调度,因此先检查是否已有存在的任务
参数:
every: 间隔数量
period: 间隔单位 (SECONDS, MINUTES, HOURS, DAYS)
task: 任务函数名(任务名)
"""
with get_session() as session:
# # 先查询是否已存在相同任务名的间隔任务
# existing_task = session.query(PeriodicTask).filter_by(
# task=task,
# discriminator='intervalschedule' # 只查询间隔任务
# ).first()
# if existing_task:
# # 如果存在,则更新celery_intervalschedule表的时间
# # 更新现有任务的调度配置
# session.execute(
# IntervalSchedule.__table__.update()
# .where(IntervalSchedule.id == existing_task.schedule_id)
# .values(every=every, period=period)
# )
# # 查询或创建一个1分钟一次的间隔调度
# schedule = session.query(IntervalSchedule).filter_by(
# every=every,
# period=period #这里的值有'days,hours,minutes,seconds'
# ).first()
# # 如果调度不存在,则创建一个
# if not schedule:
# 直接创建一个间隔调度,因为有初始化程序存在会自动删除所以不再查询
schedule = IntervalSchedule(every=every, period=period)
session.add(schedule)
session.flush()
# 创建间隔任务
interval_task = PeriodicTask(
schedule_model=schedule,
name=task,
task=task,
enabled=True # 启用任务
# 周期任务不要设置expires,因为expires参数是绝对时间,不是相对时间
)
session.add(interval_task)
# 通知调度器有新任务
session.execute(
PeriodicTaskChanged.__table__.update().values(
last_update=datetime.now()
)
)
# 创建周期任务(crontab调度)
def create_crontab_tasks():
"""创建周期任务(传入自定义的任务对象数组,使用数组批量对数据库性开销好),每个对象的属性如下:
参数:
task: 任务函数名(任务名)
minute: 分钟
"""
with get_session() as session:
# 第一个周期任务
task_name = 'echo_crontab'
# 创建第一个crontab调度
schedule = CrontabSchedule(
minute='5', # 第5分钟
# hour='8,11,14', # 8点、12点、14点
# day_of_week='*', # 每天(星期一到星期日)
# day_of_month='*', # 每月的任何一天
# month_of_year='*', # 任何月份
timezone='Asia/Shanghai', # 北京时间
)
session.add(schedule)
session.flush()
# 创建第一个周期任务
crontab_task = PeriodicTask(
schedule_model=schedule,
name=task_name,
task=task_name,
enabled=True,
# 周期任务不要设置expires,因为expires参数是绝对时间,不是相对时间
)
session.add(crontab_task)
# 第二个周期任务
task_name2 = 'tidy_celery_data'
# 创建第二个crontab调度
schedule2 = CrontabSchedule(
minute='21', # 第10分钟
hour='14,15,16', # 8点、12点、14点
timezone='Asia/Shanghai', # 北京时间
)
session.add(schedule2)
session.flush()
# 创建第二个周期任务
crontab_task2 = PeriodicTask(
schedule_model=schedule2,
name=task_name2,
task=task_name2,
enabled=True,
)
session.add(crontab_task2)
# 通知调度器有新任务
session.execute(
PeriodicTaskChanged.__table__.update().values(
last_update=datetime.now()
)
)
# 初始化调度数据库
@beat_init.connect
def init_beat(sender, **kwargs):
print("初始化调度数据库")
init_db()
# 关闭worker时,关闭数据库连接
@worker_shutdown.connect
def close_db_connection(sender, **kwargs):
print("关闭worker时,关闭数据库连接")
session_manager.session_factory(env.CELERY_DATABASE_URL).close()
tasks.py:
# 任务函数页,这里只编写需要被调度执行的函数
from datetime import datetime,timezone,timedelta # 导入时间模块
from celery_app import app,create_one_tasks,TaskOneObj,create_interval_tasks,init_db,create_crontab_tasks,get_session,ClockedSchedule,ExecutedOneTasks,PeriodicTask # 导入app和创建任务函数,任务类
# 定义一个任务函数
@app.task(name='echo')
def echo(data):
"""打印接收到的数据"""
print(f"执行echo: {data}")
# 定义一个间隔任务函数
@app.task(name='echo_interval',expires=1800)
def echo_interval():
"""打印接收到的数据"""
print(f"分钟间隔任务执行一次")
# 定义一个周期任务函数
@app.task(name='echo_crontab',expires=1800)
def echo_crontab():
"""打印接收到的数据"""
print(f"周期任务执行一次")
# 定义第二个周期任务,此任务用来定期清理数据库表
@app.task(name='tidy_celery_data',expires=1800)
def tidy_celery_data():
"""
定期清理数据库表,提升查询性能
- 删除过期的ClockedSchedule表数据
- 将超过40分钟的一次性任务迁移到历史表
"""
with get_session() as session:
# 删除所有过期超过40分钟的celery_clockedschedule任务
session.query(ClockedSchedule).filter(
ClockedSchedule.clocked_time < datetime.now() - timedelta(minutes=40)
).delete()
# 将超过40分钟的一次性任务迁移到历史表
cutoff_time = datetime.now() - timedelta(minutes=40)
# 查询需要迁移的一次性任务
old_one_off_tasks = session.query(PeriodicTask).filter(
PeriodicTask.discriminator == 'clockedschedule', # 一次性任务
PeriodicTask.date_changed < cutoff_time # 超过40分钟
).all()
# 将任务数据插入到历史表
for task in old_one_off_tasks:
# 创建历史记录
executed_task = ExecutedOneTasks(
task_name=task.name,
task=task.task,
task_args=task.args,
last_run_at=task.last_run_at if task.last_run_at else task.date_changed
)
session.add(executed_task)
# 从主表删除已迁移的任务
if old_one_off_tasks:
task_ids = [task.id for task in old_one_off_tasks]
session.query(PeriodicTask).filter(
PeriodicTask.id.in_(task_ids)
).delete(synchronize_session=False)
print(f"已执行整理数据库任务")
# 测试
if __name__ == '__main__':
# pass
# 初始化数据库,如果是第一次运行
init_db()
#创建或更新间隔任务
create_interval_tasks(2,'minutes','echo_interval')
#创建或更新周期任务(包含数据清理任务)
create_crontab_tasks()
# # 创建一次性任务对象(可以多个)
# tasks_obj = [
# TaskOneObj(task='echo',clocked_time=datetime.now(tz=timezone.utc) + timedelta(minutes=1),task_args=["一次性任务1分钟后执行2"]),
# TaskOneObj(task='echo',clocked_time=datetime.now(tz=timezone.utc) + timedelta(minutes=2),task_args=["一次性任务2分钟后执行2"]),
# ]
# # 执行创建一次任务函数
# create_one_tasks(tasks_obj)
# celery -A celery_app:app worker -P gevent -c 10
# celery -A celery_app:app beat -l info -S sqlalchemy
自定义调度器数据库+重连机制
指定了beat_scheduler就不需要在启动命令中使用 -s参数了,不然会覆盖调度器
#======顶部添加自动探活重连猴子补丁==========
from sqlalchemy_celery_beat import session as _scb_session
_orig_get_engine = _scb_session.SessionManager.get_engine
def _get_engine_with_ping(self, dburi, **kw):
kw.setdefault('pool_pre_ping', True) # 探活
kw.setdefault('pool_recycle', 3600) # 1h 强制重连
return _orig_get_engine(self, dburi, **kw)
_scb_session.SessionManager.get_engine = _get_engine_with_ping
#================猴子补丁结束==========
from sqlalchemy_celery_beat.schedulers import DatabaseScheduler #导入调度器基类
from sqlalchemy.exc import OperationalError
# 自定义调度器,避免beat进程崩溃
class SafeDatabaseScheduler(DatabaseScheduler):
"""
在父类 DatabaseScheduler 基础上增加异常捕获及有限次重试,
避免 PostgreSQL 临时不可用导致 beat 进程崩溃。
"""
_retry_delay = 5 # 单次重试间隔(秒)
_max_retry = 3 # 每个 tick 最多重试次数
def tick(self, *args, **kwargs):
"""
重写 tick:
1. 正常调用父类实现。
2. 捕获数据库 OperationalError 并循环重试。
3. 超出重试次数后返回 _retry_delay,保持 beat 存活。
"""
for attempt in range(1, self._max_retry + 1):
try:
return super().tick(*args, **kwargs)
except OperationalError as exc:
print(
"数据库连接异常: %s; 第 %s/%s 次重试,%s 秒后再试",
exc, attempt, self._max_retry, self._retry_delay
)
time.sleep(self._retry_delay)
# 多次重试仍失败,记录错误并保持 beat 继续循环
print("数据库长时间不可用,跳过本次调度周期。")
# 初始化时也设置重连
def __init__(self, *a, **kw):
for attempt in range(1, self._max_retry + 1):
try:
super().__init__(*a, **kw)
break
except OperationalError as exc:
print(f"数据库连接异常: {exc}; "
f"第 {attempt}/{self._max_retry} 次重试,"
f"{self._retry_delay} 秒后再试")
time.sleep(self._retry_delay)
return self._retry_delay
# 以下配置都写在这里
app.conf.update(
# 使用自定义调度器
{
'beat_scheduler': 'crowdpulse.celery_ser.celery_app:SafeDatabaseScheduler',
'beat_dburi': CELERY_POSTGRES_URL, # 设置beat数据库uri
'beat_schema': None, # 设置beat数据库schema
'include': ['crowdpulse.celery_ser.tasks'], # 导入任务模块
# 时区配置
'timezone':'Asia/Shanghai',
)
核心机制
- 数据库驱动:将Celery Beat的调度信息存储在关系型数据库中,替代默认的文件存储方式
- 轮询检测:Beat调度器定期扫描数据库中的任务表,根据调度规则决定何时执行任务
- 变更通知:通过
celery_periodictaskchanged表的时间戳变化来通知调度器重新加载任务 - 状态管理:跟踪任务执行状态、次数、最后执行时间等信息
调度流程
1. Beat启动 → 连接数据库 → 加载所有启用的任务
2. 轮询检查 → 判断任务是否到达执行时间
3. 发送任务 → 将任务消息发送到消息队列
4. 更新状态 → 记录执行时间、次数等信息
5. 重复循环 → 继续监控下一轮任务
数据库结构解析
1. celery_clockedschedule - 一次性定时任务调度表
存储指定时间点执行的任务调度信息,适合一次性任务
2. celery_crontabschedule - 类Cron周期性任务调度表
存储类似Linux crontab的周期性调度规则
用途:用于"每天上午9点"、"每周一下午3点"这类周期性任务
3. celery_intervalschedule - 固定间隔任务调度表
存储按固定时间间隔执行的调度规则
用途:用于"每5分钟执行一次"、"每2小时执行一次"这类间隔任务
4. celery_solarschedule - 太阳时调度表
基于日出日落时间的调度规则(较少使用)
用途:用于"日出时执行"、"日落后1小时执行"这类基于天文时间的任务
5. celery_periodictask - 核心任务定义表
存储所有定时任务的详细信息
6. celery_periodictask 核心字段解析
| 字段名 | 类型 | 说明 | 示例 |
|---|---|---|---|
name | VARCHAR(200) | 任务唯一标识名称 | "send_email_1703123456_742" |
task | VARCHAR(200) | 对应的Celery任务函数名 | "tasks.send_email" |
args | TEXT | 任务位置参数(JSON) | '["user@example.com", "Hello"]' |
kwargs | TEXT | 任务关键字参数(JSON) | '{"priority": "high"}' |
enabled | BOOLEAN | 任务是否启用 | true / false |
one_off | BOOLEAN | 是否为一次性任务 | true(执行一次后自动禁用) |
last_run_at | TIMESTAMP | 最后执行时间 | 2024-01-01 10:30:00+00 |
total_run_count | INTEGER | 累计执行次数 | 15 |
schedule_type | VARCHAR(50) | 调度类型标识 | "clocked", "crontab", "interval" |
schedule_id | INTEGER | 关联调度表的ID | 指向具体调度表的记录ID |
one_off | BOOLEAN | 是否为一次性任务 | true |
start_time | TIMESTAMP | 任务开始生效时间 | 2024-01-01 09:00:00+00 |
expires | TIMESTAMP | 任务超时绝对时间 | 2024-01-01 09:00:00+00 |
expire_seconds | INTEGER | 超时时间。以秒为单位的相对时间 | 6 |
priority | INTEGER | 任务优先级(0-9) | 6 |
7. celery_periodictaskchanged - 任务变更通知表
用于通知Beat调度器重新加载任务配置
重要特性:
- 表中只有一条记录
- 每次任务变更时更新
last_update字段 - Beat调度器通过监控此字段变化来判断是否需要重新加载任务
任务生命周期
对于一次性任务建议,只设置clocked_time即可,到点就提交
clocked_time 任务发送到队列时间
start_time 任务生效时间,这个一般可以不设置
一次性任务 (one_off=True)
创建 → enabled=True, total_run_count=0
执行 → enabled=False, total_run_count=1, last_run_at=执行时间
完成 → 保留记录用于审计(可手动清理)
周期性任务 (one_off=False)
创建 → enabled=True, total_run_count=0
执行 → enabled=True, total_run_count++, last_run_at=执行时间
循环 → 继续等待下次执行时间
性能优化建议
1. 数据库索引
-- 推荐创建的索引
CREATE INDEX idx_periodictask_enabled ON celery_periodictask(enabled);
CREATE INDEX idx_periodictask_next_run ON celery_periodictask(enabled, last_run_at);
CREATE INDEX idx_periodictask_one_off ON celery_periodictask(one_off, enabled);
2. 定期清理
- 一次性任务清理:删除已执行的
one_off=true且enabled=false的记录 - 调度记录清理:删除不再被
periodictask引用的调度记录 - 日志清理:定期清理过期的执行日志
启动
启动celery消费者
celery -A celery_app:app worker -P gevent -c 10
启动celery调度器(注意这里使用了-S,指定了使用sqlalchemy_celery_beat(简写sqlalchemy)作为调度器)
celery -A celery_app:app beat -l info -S sqlalchemy
go版本的类似任务队列方案
PostgreSQL + app_schedules + River