Celery SQLAlchemy 数据库结果后端(celery.backends.database)源码与配置深度解析
【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery
本篇技术指南围绕 Celery 分布式任务队列中的 SQLAlchemy 数据库结果后端(celery.backends.database)展开,系统讲解其模块结构、ORM 数据模型、会话管理机制、全部相关配置项(database_url、database_engine_options、database_short_lived_sessions、database_table_schemas/database_table_names、database_engine_callback、database_create_tables_at_setup等)以及任务结果读写、分组结果保存、过期清理与失败重试的完整实现链路。读者读完后,将能依据仓库源码与官方配置文档,独立完成数据库结果后端的选型、配置、表结构定制、升级迁移与故障排查。
模块定位:SQLAlchemy 结果存储后端
celery.backends.database是 Celery 中以 SQLAlchemy ORM 为底层实现的任务结果存储后端。它把每个任务的执行状态(PENDING/STARTED/SUCCESS/FAILURE等)、返回值、异常 traceback、执行时间以及任务子结果(children)持久化到关系型数据库中,并支持保存与恢复 Group(任务组)级联结果。
该模块的对外唯一入口类为DatabaseBackend,声明于 celery/backends/database/init.py,继承自celery.backends.base.BaseBackend。模块整体由三个文件协作组成:
- celery/backends/database/init.py:
DatabaseBackend后端实现与可重试数据库错误集合; - celery/backends/database/models.py:
Task、TaskExtended、TaskSet三个 ORM 模型; - celery/backends/database/session.py:
SessionManager会话/引擎管理与建表(含缺失列自动迁移)逻辑。
数据模型:两张核心表与三种模型类
从 celery/backends/database/models.py 可以看到,后端自动维护两张表,所有模型均基于ResultModelBase(由sqlalchemy.orm.declarative_base()生成,定义于 celery/backends/database/session.py)。
Task:celery_taskmeta表
任务级结果与状态表,表名默认celery_taskmeta,字段如下:
| 字段 | 列类型 | 说明 |
|---|---|---|
id | 整型主键(MSSQL 下自动切换为BigInteger变体),自增 | 内部自增主键,并绑定task_id_sequence序列 |
task_id | String(155),唯一索引 | Celery 任务 UUID |
status | String(50),默认PENDING | 任务状态,取值对应 celery/states.py 中的状态集合 |
result | LargeBinary,可空 | 序列化后的任务结果(自 Celery 5.7 起遵循result_serializer) |
date_done | DateTime,可空,带索引 | 完成时间;默认值使用_get_utc_now()这一可调用对象,确保在 INSERT/UPDATE 时实时求值(而非模块导入时固化) |
traceback | Text,可空 | 失败任务的异常回溯 |
children | LargeBinary,可空 | 子任务结果(5.7 起支持),同样以配置的序列化器编码 |
值得注意的是id列使用DialectSpecificInteger = sa.Integer().with_variant(sa.BigInteger, 'mssql'),即 SQL Server 方言下自动使用大整数,这是对 MSSQL 主键溢出场景的方言适配。
TaskExtended:扩展字段模型
当开启result_extended = True时,后端切换使用TaskExtended(同样映射到celery_taskmeta表,通过extend_existing=True复用表定义)。它在Task基础上追加name(任务名)、args(序列化参数)、kwargs(序列化关键字参数)、worker(执行 worker)、retries(重试次数)、queue(队列名)六个可空列,全部为String/LargeBinary/Integer类型。切换逻辑位于DatabaseBackend.__init__:if self.extended_result: self.task_cls = TaskExtended。
TaskSet:celery_tasksetmeta表
Group 级结果表,默认表名celery_tasksetmeta,字段包括id(自增主键)、taskset_id(String(155)唯一)、result(LargeBinary,存组结果序列化值)、date_done(带索引)。Group 结果的写入与读取由_save_group/_restore_group负责,读取时通过result_from_tuple还原为GroupResult对象。
表名与 Schema 的动态配置
Task.configure()与TaskSet.configure()类方法支持在运行时修改表所属 schema 与表名。DatabaseBackend.__init__中据此从配置读取:
schemas = conf.database_table_schemas or {} tablenames = conf.database_table_names or {} self.task_cls.configure(schema=schemas.get('task'), name=tablenames.get('task')) self.taskset_cls.configure(schema=schemas.get('group'), name=tablenames.get('group'))对应配置示例(详见 docs/userguide/configuration.rst):
# 自定义表所属 schema(PostgreSQL 等支持 schema 的数据库) database_table_schemas = {'task': 'celery', 'group': 'celery'} # 自定义表名 database_table_names = {'task': 'myapp_taskmeta', 'group': 'myapp_groupmeta'}会话管理:SessionManager 与短生命周期会话
celery/backends/database/session.py 中的SessionManager负责 SQLAlchemy Engine 与 Session 的创建、缓存、fork 安全与失效清理:
- 引擎创建:
get_engine()在非 fork 场景下使用NullPool(不维护持久连接池),并过滤掉max_overflow、echo_pool等 NullPool 不支持的参数;fork 后(如 prefork worker 的子进程)则复用已缓存的引擎,避免重复创建。 - 注册 after_fork 钩子:构造时通过 kombu 的
register_after_fork注册_after_fork_cleanup_session,fork 后将forked标志置真,保证子进程不会误用父进程的连接。 - 失效清理:
invalidate(dburi)弹出并dispose()对应引擎,用于在可重试数据库错误后重置连接状态。 - 建表与自动迁移:
prepare_models()调用ResultModelBase.metadata.create_all(engine),并对DatabaseError采用指数退避重试(最多PREPARE_MODELS_MAX_RETRIES = 10次);随后_migrate_missing_columns()检查存量表是否缺少可空列(如 5.7 新增的children列),缺失时自动执行ALTER TABLE ... ADD COLUMN,若当前数据库用户无 DDL 权限则记录 warning 并优雅降级。
DatabaseBackend.ResultSession()是对外的会话工厂,组合了dburi、short_lived_sessions与合并后的引擎选项。short_lived_sessions为真时每次操作都新建 sessionmaker 绑定(见create_session),官方文档明确指出:默认关闭,开启后会显著降低大批量任务场景的性能,但能解决低流量 worker 因连接闲置过期引发的(OperationalError) (2006, 'MySQL server has gone away')类错误。
每个数据库操作都包裹在session_cleanup上下文管理器(celery/backends/database/init.py)中:异常时rollback()后重新抛出,finally中总是close()会话,避免连接泄漏。
完整配置指南
所有与数据库后端相关的配置项默认值集中定义在 celery/app/defaults.py 的databaseNamespace 中,官方说明见 docs/userguide/configuration.rst。以下逐项展开。
1. 结果后端地址:result_backend与database_url
使用数据库后端必须配置带db+前缀的连接 URL:
result_backend = 'db+scheme://user:password@host:port/dbname'官方文档给出的四类典型示例:
# sqlite(文件型数据库) result_backend = 'db+sqlite:///results.sqlite' # mysql result_backend = 'db+mysql://scott:tiger@localhost/foo' # postgresql result_backend = 'db+postgresql://scott:tiger@localhost/mydatabase' # oracle result_backend = 'db+oracle://scott:tiger@127.0.0.1:1521/sidname'在源码层面,URL 的解析优先级为url or dburi or conf.database_url(见DatabaseBackend.__init__),其中url参数由celery.app.backends.by_url按db+前缀路由传入。若三者均缺失,会抛出ImproperlyConfigured,提示"Missing connection string!"。database_url配置项的默认值类型为Option(old={'celery_result_dburi'}),即兼容旧版CELERY_RESULT_DBURI环境变量命名。
2. 引擎选项:database_engine_options
默认值为{'pool_pre_ping': True, 'pool_recycle': 3600}(自 5.7 起从空字典调整而来),用于改善连接健康:pool_pre_ping在取连接时先探测有效性,pool_recycle强制连接 3600 秒后回收,可有效避免 MySQL 等数据库空闲断连导致的 stale connection 错误。
# 开启 SQLAlchemy 详细 SQL 日志 app.conf.database_engine_options = {'echo': True} # 显式关闭默认的连接健康选项 app.conf.database_engine_options = {'pool_pre_ping': False, 'pool_recycle': None}源码中引擎选项的合并规则为:conf.database_engine_options(配置默认值)为底,构造器传入的engine_options覆盖其上(dict(conf.database_engine_options or {}, **(engine_options or {})))。
3. 会话策略:database_short_lived_sessions
默认False。启用后每次操作使用全新会话,适合低流量、易遇陈旧连接的场景,但会显著影响高吞吐 worker 性能(官方文档明确警示)。
4. 建表时机:database_create_tables_at_setup
True(默认,5.5 起):后端初始化(DatabaseBackend.__init__)时立即调用_create_tables()建表;False:延迟到第一个任务执行、首次创建会话时由prepare_models()惰性建表(等价于 5.5 之前的行为)。
5. 引擎回调:database_engine_callback
5.7 新增,默认None。可以是可调用对象或点分导入路径字符串,引擎创建后立即被调用,用于注册do_connect等事件监听器,实现 JWT 令牌注入、IAM 鉴权等按连接认证需求:
from sqlalchemy import event def register_do_connect(engine): @event.listens_for(engine, 'do_connect') def on_connect(dialect, conn_rec, cargs, cparams): cparams['password'] = get_auth_token() app.conf.database_engine_callback = register_do_connect # 或以字符串形式指定 app.conf.database_engine_callback = 'myapp.db:register_do_connect'字符串形式在源码中通过celery.utils.imports.symbol_by_name解析;若非可调用对象则抛出ImproperlyConfigured。
6. 重试策略:result_backend_always_retry与result_backend_max_retries
为保证向后兼容,数据库后端覆盖了BaseBackend的默认值(celery/backends/base.py 中分别为False与inf):数据库后端默认result_backend_always_retry = True、result_backend_max_retries = 3(见DatabaseBackend.__init__中注释:此前版本使用自定义@retry装饰器固定重试 3 次)。
# 关闭自动重试 result_backend_always_retry = False # 提升重试上限 result_backend_always_retry = True result_backend_max_retries = 10哪些异常值得重试由exception_safe_to_retry()判定,RETRYABLE_DB_ERRORS包含DatabaseError、InterfaceError、InvalidRequestError、StaleDataError。每次可重试错误触发后,on_backend_retryable_error()会调用session_manager.invalidate(self.url)丢弃失效引擎与会话,避免在坏连接上反复重试。
7. 序列化:result_serializer与旧数据兼容
自 Celery 5.7 起,result与children列内容遵循result_serializer配置(官方文档专门为此新增说明)。为兼容 5.7 之前"无论配置什么序列化器都写 pickle"的历史数据,_decode_stored_result()在解码失败时检查载荷首字节是否为 pickle 协议 2+ 标记\x80,是则回退pickle.loads()并输出 warning,否则原样抛出解码错误,杜绝盲目反序列化任意字节的安全隐患(对应 issue celery/celery#3025)。
8. 扩展结果:result_extended
设为True后后端改用TaskExtended模型,额外持久化name、args、kwargs、worker、retries、queue六个字段;读取时_get_task_meta_for会对args/kwargs/children执行decode还原。
核心操作流程与调用链
以 celery/backends/database/init.py 的实现为准,后端六大核心操作如下:
存储任务结果:_store_result
新建会话 →_query_task按task_id查表(命中DatabaseError且错误信息含children时,回退为defer(children)懒加载查询以绕过列缺失问题)→ 不存在则Task(task_id)并session.add/flush→_update_result依据_get_result_meta生成的元数据逐列setattr(显式排除主键id、task_id与单独处理的children)→session.commit()。
查询任务元数据:_get_task_meta_for
查询task_id;不存在时构造一个status = PENDING、result = None的占位Task,保证未执行任务的元数据查询也能返回合法的 pending 结构。随后task.to_dict()转字典,result经_decode_stored_result解码,args/kwargs/children分别decode,最后交由meta_from_decoded归一化。
分组结果:_save_group/_restore_group/_delete_group
组结果以ensure_bytes(self.encode(self.prepare_value(result)))编码后写入TaskSet;读取时先解码再经result_from_tuple(value, self.app)还原为GroupResult;删除时按taskset_id执行 DELETE。
忘记结果:_forget
按task_id删除celery_taskmeta中的记录行。
过期清理:cleanup/_cleanup
计算now - expires作为截止时间,一次性删除date_done早于该时间的任务与组记录。这正是内置周期任务celery.backend_cleanup的后端实现。因此官方文档建议为date_done添加数据库索引以提升大表清理性能,升级到 5.7 后可通过 Alembic 迁移或手工 SQL 补建索引:
CREATE INDEX ix_celery_taskmeta_date_done ON celery_taskmeta (date_done); CREATE INDEX ix_celery_tasksetmeta_date_done ON celery_tasksetmeta (date_done);结果存在性检查:task_result_exists(5.7 新增)
通过_ensure_retryable包装的_task_result_exists实现,返回布尔值,用于避免对已完成任务重复执行等场景。
可重试机制与测试佐证
数据库后端的重试实现基于BaseBackend._ensure_retryable包装层(always_retry/max_retries/exception_safe_to_retry/on_backend_retryable_error四者的协作)。单元测试 t/unit/backends/test_database.py 对上述行为做了系统验证,包括:
test_store_result_retries_on_database_error、test_get_task_meta_retries_on_database_error、test_save_group_retries_on_database_error、test_forget_retries_on_database_error、test_cleanup_retries_on_database_error:各核心操作在DatabaseError下自动重试;test_retries_respect_max_retries_config:重试次数遵循result_backend_max_retries配置;test_non_retryable_exceptions_propagate_immediately:非可重试异常立即上抛;test_engine_options_include_pool_health_defaults/test_engine_options_explicit_values_override_defaults:验证pool_pre_ping、pool_recycle默认值与显式覆盖;test_table_schema_config/test_table_name_config:验证database_table_schemas/database_table_names生效;test_result_respects_configured_serializer/test_result_decode_falls_back_to_pickle_for_legacy_rows:验证 5.7 序列化行为与旧 pickle 数据回退;test_migrate_missing_columns/test_migrate_missing_columns_failure_logs_warning_and_degrades_gracefully:验证children列自动迁移及无权限时的优雅降级;test_missing_task_id_is_PENDING:验证未找到任务时返回 pending 占位。
升级与运维注意事项
综合官方配置文档(docs/userguide/configuration.rst)与源码,从旧版本升级或运维数据库结果后端时有四点关键事项:
- 序列化行为变更(5.7):新写入的结果遵循
result_serializer,旧 pickle 数据仍可自动回退读取,无需迁移 schema; children列迁移(5.7):后端启动时会自动为存量表补加可空children列;若数据库账号无 DDL 权限,需手动执行ALTER TABLE celery_taskmeta ADD COLUMN children BLOB;(SQLite/MySQL)或ADD COLUMN children BYTEA;(PostgreSQL);date_done索引(5.7 建议):create_all()不会修改已存在的表,大表上请按上文 SQL 补建索引以加速celery.backend_cleanup清理;- 连接健康默认值(5.7):
database_engine_options默认引入pool_pre_ping=True与pool_recycle=3600,若自定义该配置会整体覆盖默认值,需要显式保留这两个键。
小结
celery.backends.database通过 DatabaseBackend、ORM 模型 与 SessionManager 三部分,把任务结果可靠地落到任意 SQLAlchemy 支持的数据库中。它既保留了历史版本的惰性建表与 pickle 兼容读取,又在 5.7 起补齐了连接健康检查、可重试错误处理、children持久化、自动列迁移与引擎回调等生产级能力。结合 t/unit/backends/test_database.py 中的测试用例与官方配置文档,开发者可以放心地在 MySQL、PostgreSQL、SQLite、Oracle 等数据库上将其作为 RPC/Redis 之外的持久化结果后端使用。
【免费下载链接】celeryDistributed Task Queue (development branch)项目地址: https://gitcode.com/gh_mirrors/ce/celery
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考