Airflow Task State Store 数据如何定期清理与设置保留期?
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
如果你在用 Airflow 的 task state store(3.3 版本引入的持久化 key/value 存储,用于保存外部任务 ID、检查点等),会遇到一个运维问题:这些行不会自动过期删除。Airflow 不会按计划清除 task state store 的行,清理(garbage collection)必须由用户显式通过 CLI 触发。本文说明如何为这些行设置保留期,以及如何把airflow state-store clean纳入定期维护流程。内容适用 Airflow 3.3 及以上版本,且仅对默认的 metastore backend(数据存在 Airflow 元数据库中)有效。
先弄清哪些行会被清理
清理命令只作用于task state store的行,asset state store 的行永远不会被该命令触碰——asset 行只在 asset 被停用时由 orphan sweep 删除。
一个 task state store 行只有expires_at时间戳已过期才会被删除。expires_at是在 worker 上写入时计算的,规则如下(来源:清理文档 与 task state store 概念文档):
- 写入时显式指定
retention=timedelta(...)的 key,在该时长后过期; - 写入时
retention=None(默认)的 key,按[state_store] default_retention_days计算过期时间; retention=NEVER_EXPIRE的 key 存储为expires_at = NULL并带有永久标记,无论配置如何都不会被该命令删除;- 若
default_retention_days = 0,未显式指定 retention 的 key 也没有过期时间,同样被跳过。
只有expires_at非空且已过去的行会被删除。
在 airflow.cfg 中设置保留期
所有相关配置都在airflow.cfg的[state_store]段落里。注意段落名是[state_store],不是[task_state_store](配置参考 中特别强调了这一点)。
核心配置项及默认值:
# airflow.cfg [state_store] # 无显式 retention 的 key 在写入 N 天后过期,默认 30;设为 0 完全禁用基于时间的清理 default_retention_days = 30 # 清理时每个批次删除的行数,默认 0(单条语句删完) state_cleanup_batch_size = 0 # 任务实例进入 success 状态时自动删除其全部 task state store key,默认 False clear_on_success = False使用上的几个关键点:
default_retention_days只影响 task state store,不影响 asset state store 行。- 修改该配置不会作用于已经写入的行——
expires_at在写入时已经算好,所以调整保留期只对之后新写入的 key 生效。 - 代码侧写入时也可以按 key 精细控制:
task_state_store.set("key", "val", retention=timedelta(days=7))。文档特别提醒retention只接受datetime.timedelta,传整数会抛TypeError;需要永不过期时用from airflow.sdk import NEVER_EXPIRE。 clear_on_success = True是一种可选的补充路径:任务成功即删行,不依赖保留期。它只清理 task state store,对 asset store 无效。如果你不需要成功后的可观测性(例如从 UI/REST API 回看提交的 job ID),可以打开它以在不等待保留期的情况下自动清理。
运行清理命令
清理命令是:
airflow state-store clean它会读取[state_store] default_retention_days和[state_store] state_cleanup_batch_size,然后删除所有已过期的行。
先做 dry run 预览。加--dry-run只列出将被删除的行,不做任何删除:
airflow state-store clean --dry-run输出按 dag、run、task、map index 和 key 分组列出每一行会被删除的记录(格式见命令实现 state_store_command.py,示例输出):
Would delete 2 task state store row(s): Dag 'paginated_ingest', run 'manual__2026-09-01T00:00:00+00:00', task 'ingest_pages', map_index -1, key 'last_page' Dag 'row_ingest', run 'scheduled__2026-09-01T00:00:00+00:00', task 'ingest_rows', map_index -1, key 'progress'没有可删行时输出Nothing to delete.。建议先 dry run 确认范围,再执行正式清理。
大表上设置批次大小。默认state_cleanup_batch_size = 0时,所有符合条件的行在单条语句中删除。如果你的task_state_store表很大,设置批次大小可以降低每个事务的锁持有时长:
# airflow.cfg [state_store] state_cleanup_batch_size = 10000命令会按每批 10,000 行删除、每批提交一次,直到没有符合条件的行。
定期执行的频率
文档没有内置定时机制,需要把该命令纳入你自己的周期性维护任务(例如外部调度器按计划调用 CLI)。频率选择依据写入量:
- 多数环境每周清理一次即可;
- 对每次任务执行都写入 task state store 的高吞吐管道,建议提高清理频率,以控制
task_state_store表的大小。
限制与不适用情况
- 自定义 backend 会被跳过。若
[state_store] backend指向非默认实现,清理命令会打印一条消息并正常退出,不删除任何内容。如果自定义 backend 需要保留期逻辑,要在BaseStoreBackend.cleanup()中自行实现并调用。 NEVER_EXPIRE的 key 永不清理。用它们保存外部任务 ID 的持久化执行场景不受此命令影响。- asset state store 不在此命令范围内,它由 asset 停用时的 orphan sweep 处理。
- 调整保留期不追溯生效,只对新写入的 key 起作用。
更多背景(后端语义、自定义 backend 实现方式)见 Task and Asset State Store 配置文档。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考