celery工程实践


使用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:

python
# 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:

python
# 任务函数页,这里只编写需要被调度执行的函数
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参数了,不然会覆盖调度器

python
#======顶部添加自动探活重连猴子补丁==========
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',
)

核心机制

  1. 数据库驱动:将Celery Beat的调度信息存储在关系型数据库中,替代默认的文件存储方式
  2. 轮询检测:Beat调度器定期扫描数据库中的任务表,根据调度规则决定何时执行任务
  3. 变更通知:通过 celery_periodictaskchanged 表的时间戳变化来通知调度器重新加载任务
  4. 状态管理:跟踪任务执行状态、次数、最后执行时间等信息

调度流程

markdown
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 核心字段解析

字段名类型说明示例
nameVARCHAR(200)任务唯一标识名称"send_email_1703123456_742"
taskVARCHAR(200)对应的Celery任务函数名"tasks.send_email"
argsTEXT任务位置参数(JSON)'["user@example.com", "Hello"]'
kwargsTEXT任务关键字参数(JSON)'{"priority": "high"}'
enabledBOOLEAN任务是否启用true / false
one_offBOOLEAN是否为一次性任务true(执行一次后自动禁用)
last_run_atTIMESTAMP最后执行时间2024-01-01 10:30:00+00
total_run_countINTEGER累计执行次数15
schedule_typeVARCHAR(50)调度类型标识"clocked", "crontab", "interval"
schedule_idINTEGER关联调度表的ID指向具体调度表的记录ID
one_offBOOLEAN是否为一次性任务true
start_timeTIMESTAMP任务开始生效时间2024-01-01 09:00:00+00
expiresTIMESTAMP任务超时绝对时间2024-01-01 09:00:00+00
expire_secondsINTEGER超时时间。以秒为单位的相对时间6
priorityINTEGER任务优先级(0-9)6

7. celery_periodictaskchanged - 任务变更通知表

用于通知Beat调度器重新加载任务配置
重要特性

  • 表中只有一条记录
  • 每次任务变更时更新 last_update 字段
  • Beat调度器通过监控此字段变化来判断是否需要重新加载任务

任务生命周期

对于一次性任务建议,只设置clocked_time即可,到点就提交
clocked_time 任务发送到队列时间
start_time 任务生效时间,这个一般可以不设置

一次性任务 (one_off=True)

ini
创建 → enabled=True, total_run_count=0
执行 → enabled=False, total_run_count=1, last_run_at=执行时间
完成 → 保留记录用于审计(可手动清理)

周期性任务 (one_off=False)

ini
创建 → enabled=True, total_run_count=0
执行 → enabled=True, total_run_count++, last_run_at=执行时间
循环 → 继续等待下次执行时间

性能优化建议

1. 数据库索引

sql
-- 推荐创建的索引
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=trueenabled=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