PostHog 异步迁移(Async Migrations)完全指南:从编写、运行到回滚
【免费下载链接】posthog:hedgehog: PostHog is the leading platform for building self-driving products. Our developer tools – AI observability, analytics, session replay, flags, experiments, error tracking, logs, and more – capture all the context agents need to diagnose problems, uncover opportunities, and ship fixes. Steer it all from Slack, web, desktop, or the MCP.项目地址: https://gitcode.com/GitHub_Trending/po/posthog
导读
异步迁移(Async Migrations)是 PostHog 在 Django/EE 同步迁移之外提供的一套后台数据迁移机制,专门用于处理无法在服务启动阶段同步完成的重量级 ClickHouse/Postgres 数据变更(例如替换 ClickHouse 表引擎、回填数据、重建物化视图)。本文将基于仓库内 异步迁移工程文档 以及posthog/async_migrations/目录下的源码实现,完整讲解迁移文件的编写规范、底层工作流与架构、启动前置检查、健康检查、停止与回滚机制、相关配置项,以及整个代码库的结构,帮助你快速上手编写一个符合 PostHog 规范、可安全回滚的异步迁移。
异步迁移的定位与适用范围
在开始编写之前,先明确异步迁移的边界。根据文档及 definition.py 的实现,异步迁移的初始设计只面向数据迁移(data migrations),其核心假设是:迁移是帮助用户从旧状态过渡到新默认状态的机制。
典型场景如文档中提到的案例:当 PostHog 将 ClickHouse 的person_distinct_id表迁移到CollapsingMergeTree引擎时,代码库会同步更新建表 SQL,并编写一个异步迁移帮助仍在使用旧 schema 的用户完成升级。而在变更之后全新部署的实例,其默认建表 SQL 已经是新 schema,无需再跑迁移——这正是is_required函数存在的意义:通过检查实例当前状态决定迁移是否需要执行,从而避免在全新实例上重复执行无用迁移。
因此在编写异步迁移时,一个关键准则是:写一个合理的is_required函数,判断当前实例是否真的需要这个迁移。新部署的 PostHog 实例会先按顺序执行所有 EE 迁移,再按顺序执行所有异步迁移,此时如果代码库中已包含更新的默认 schema,异步迁移应当被跳过。
编写一个异步迁移
文件位置与命名规范
异步迁移文件应放在posthog/async_migrations/migrations/目录下,命名沿用 Django 与 EE 迁移的规范,例如0005_update_events_schema。仓库中实际存在从0001_events_sample_by到0010_move_old_partitions的系列迁移(见 posthog/async_migrations/migrations/),并提供了可直接参考的 示例迁移 与测试示例 test_migration.py、test_with_rollback_exception.py。
核心组成:AsyncMigrationDefinition
每个迁移文件必须导出一个继承自AsyncMigrationDefinition的Migration类,通过类属性声明迁移的元信息。以 examples/example.py 为例:
class Migration(AsyncMigrationDefinition): description = "An example async migration." posthog_min_version = "1.29.0" posthog_max_version = "1.30.0" service_version_requirements = [ ServiceVersionRequirement(service="clickhouse", supported_version=">=21.6.0,<21.7.0") ]参照 definition.py 中的类定义,可配置的元信息包括:
| 属性 | 类型 | 说明 |
|---|---|---|
description | str | 迁移用途说明,展示给自托管用户 |
posthog_min_version | str | 迁移可运行的最低 PostHog 版本(默认0.0.0) |
posthog_max_version | str | 迁移必须完成的最新 PostHog 版本(默认10000.0.0) |
service_version_requirements | list | 迁移所依赖服务(ClickHouse/Postgres/Redis)的版本范围 |
depends_on | Optional[str] | 本迁移依赖的其它异步迁移名称 |
parameters | dict | 可选参数,形如{name: (默认值, 描述, 校验函数)},会在管理界面启动迁移时展示 |
其中service_version_requirements使用 version_requirement.py 中的ServiceVersionRequirement类,其accepted_services限定为("clickhouse", "postgresql", "redis"),supported_version使用semantic_version库的SimpleSpec语法(如">=21.6.0,<21.7.0")。版本校验通过实时查询各服务版本实现:ClickHouse 执行SELECT version()、Postgres 执行SHOW server_version、Redis 读取INFO中的redis_version。
定义操作:AsyncMigrationOperation 与 AsyncMigrationOperationSQL
迁移的实际工作由operations属性定义的操作列表承载,类定义见 definition.py:
AsyncMigrationOperation(fn, rollback_fn):通用操作,fn接收一个query_id字符串(用于在 ClickHouse 中标记查询便于排查),rollback_fn默认是空操作。注意rollback_fn会被同步执行,不应是耗时操作;显式传入None会导致回滚失败。AsyncMigrationOperationSQL:SQL 操作,支持以下参数:sql:要执行的 SQL 语句;sql_settings:传给执行客户端的 settings(默认使用max_execution_time超时);rollback:回滚 SQL(None表示跳过回滚);rollback_settings:回滚语句的 settings;database:目标数据库,AnalyticsDBMS.CLICKHOUSE(默认)或AnalyticsDBMS.POSTGRES;timeout_seconds:超时时间,默认取ASYNC_MIGRATIONS_DEFAULT_TIMEOUT_SECONDS(见下节配置);per_shard:是否逐分片执行。
ClickHouse 操作在 utils.py 的execute_op_clickhouse中执行,会通过tag_queries(kind="async_migration", id=query_id)标记查询,并在失败时抛出带 SQL 与 query_id 的异常;per_shard=True时由execute_on_each_shard逐分片执行(受CLICKHOUSE_ALLOW_PER_SHARD_EXECUTION开关控制)。Postgres 操作由execute_op_postgres执行,同样会在语句前附加/* query_id */注释。
一个真实的 SQL 操作示例(来自 example.py,把person_distinct_id数据写入临时表并定义回滚为删除临时表):
AsyncMigrationOperationSQL( database=AnalyticsDBMS.CLICKHOUSE, sql=f""" INSERT INTO {TEMPORARY_TABLE_NAME} (distinct_id, person_id, team_id, _sign, _timestamp, _offset) SELECT distinct_id, person_id, team_id, if(is_deleted==0, 1, -1) as _sign, _timestamp, _offset FROM {PERSONS_DISTINCT_ID_TABLE} """, rollback=f"DROP TABLE IF EXISTS {TEMPORARY_TABLE_NAME}", ),生命周期钩子:is_required / precheck / healthcheck / progress
AsyncMigrationDefinition提供四个可在子类中覆写的钩子:
is_required() -> bool:启动前判断实例是否需要本迁移。文档给出的范例是检查目标表是否已存在,若已存在则跳过:
def is_required(self): result = sync_execute("SELECT count(*) FROM system.tables WHERE database='posthog' AND name='table_x_new'") return result[0][0] == 0is_required也可以结合表结构判断,例如检查SHOW CREATE TABLE的输出(参考 example.py 中通过判断引擎是否为ReplacingMergeTree来决定是否迁移)。在 runner.py 的start_async_migration中,若is_required()返回False,迁移会直接被标记为CompletedSuccessfully(跳过但不报错)。
precheck() -> tuple[bool, Optional[str]]:启动前运行的前置检查,返回(是否通过, 失败原因),默认(True, None)。healthcheck() -> tuple[bool, Optional[str]]:迁移执行期间周期性运行的健康检查。示例中的实现会检查 ClickHouse 磁盘剩余空间,不足 1/3 时返回失败:
def healthcheck(self): result = sync_execute("SELECT total_space, free_space FROM system.disks") total_space = result[0][0] free_space = result[0][1] if free_space > total_space / 3: return (True, None) else: return (False, "Upgrade your ClickHouse storage.")progress(migration_instance) -> int:返回 0–100 的进度百分比。默认实现为100 * current_operation_index / len(operations);也可按实际数据处理量计算,如示例中按"已迁移行数 / 总行数"计算进度。
此外,get_parameter(parameter_name)与migration_instance()帮助迁移在运行期读取用户通过界面传入的参数或当前数据库记录。
工作流与架构
服务启动时的 Setup
文档描述的 Setup 流程在 setup.py 的setup_async_migrations中实现,共五步:
- 导入所有迁移定义:
import_submodules(ASYNC_MIGRATIONS_MODULE_PATH)导入posthog.async_migrations.migrations下所有模块,并构建ALL_ASYNC_MIGRATIONS内存字典; - 填充依赖图与内存记录:
_set_up_dependency_constants遍历所有迁移的depends_on,构建ASYNC_MIGRATION_TO_DEPENDENCY与反向映射DEPENDENCY_TO_ASYNC_MIGRATION;同时校验只有一个迁移没有依赖(即链式依赖的起点),否则抛出ImproperlyConfigured; - 为每个迁移创建数据库记录:
setup_model通过AsyncMigration.objects.get_or_create在 Postgres 中落库,并同步description、posthog_min_version、posthog_max_version; - 检查本版本必需迁移是否全部完成:若存在未完成且
posthog_max_version低于当前FROZEN_POSTHOG_VERSION、且is_required()为真的迁移,则抛出ImproperlyConfigured,阻止服务启动(此行为受ASYNC_MIGRATIONS_BLOCK_UPGRADE控制); - 自动触发迁移:若实例设置
AUTO_START_ASYNC_MIGRATIONS开启且存在未完成的迁移,则从链首迁移开始尝试运行,完成后通过依赖链自动接力下一个迁移。
迁移的运行流程与前置检查
迁移触发后,trigger_migration(见 utils.py)会向 Celery 派发run_async_migration任务。首次启动走start_async_migration(runner.py),按文档列出的检查顺序依次验证:
- 并发数未超限:
MAX_CONCURRENT_ASYNC_MIGRATIONS = 1,即全实例同时只允许一个迁移运行(get_all_running_async_migrations().count() >= 1即拒绝); - PostHog 版本兼容:
FROZEN_POSTHOG_VERSION处于[posthog_min_version, posthog_max_version]区间(可用ASYNC_MIGRATIONS_IGNORE_POSTHOG_VERSION或ignore_posthog_version绕过); - 迁移未在运行:状态必须是
Starting(UI 触发)或NotStarted(API 触发); is_required通过:不通过则直接标记完成并返回成功;- 服务版本要求满足:逐条校验
service_version_requirements,任一不满足则以FailedAtStartup状态记录错误; - 依赖已完成:
is_migration_dependency_fulfilled要求depends_on指定的迁移状态为CompletedSuccessfully; precheck与healthcheck通过。
全部通过后mark_async_migration_as_running通过select_for_update原子地把状态置为Running(若状态已变化则放弃启动),随后run_async_migration_operations循环执行operations中的每个操作:每成功执行一个操作就递增current_operation_index、记录current_query_id并更新进度;current_operation_index超过操作总数时调用complete_migration将状态置为CompletedSuccessfully(可选发送完成邮件)。
周期性健康检查
文档指出每 30 分钟会有一个 Celery 任务执行健康检查,对应 tasks/async_migrations.py 中的check_async_migration_health。该任务负责:
- 检测 worker 崩溃:通过
AsyncResult(celery_task_id).state与app.control.inspect().active()对比,若任务 ID 不在活跃任务列表中,说明 worker 已崩溃。此时若ASYNC_MIGRATIONS_AUTO_CONTINUE开启则用fresh_start=False重新触发迁移继续执行(利用current_operation_index断点续跑);否则记录错误并触发回滚; - 周期性 healthcheck:调用迁移定义的
healthcheck(),失败则force_stop_migration强制停止并回滚; - 更新进度:调用
update_migration_progress刷新 UI 展示的进度(进度检查失败不打断迁移)。
Celery 任务本身通过@shared_task(track_started=True, ignore_result=False, max_retries=0)定义,注释中特别说明这会占用整个 worker,文档也提醒可考虑在迁移期间扩容 Celery。
补充说明:文档提到"Async migrations can also be run synchronously (i.e. not in Celery) using the async migrations CLI (WIP) or the Django shell"。从源码看,
fresh_start=False的续跑路径可直接调用run_async_migration_operations(migration_name),这与 Django shell 中手动驱动迁移的机制一致——手动调用可绕过 Celery 派发,但需自行确保并发与状态约束。
停止、回滚与错误处理
- 停止:可以从异步迁移管理页面操作,或通过 Celery app control 终止执行任务。源码中
force_stop_migration(utils.py)会先处理Starting状态(halt_starting_migration原子地置为RolledBack以阻止启动),随后app.control.revoke(celery_task_id, terminate=True)直接杀掉执行进程,最后记录错误并按需回滚。源码注释坦诚指出:terminate存在任务已完成才被杀掉的竞态窗口,但因为迁移任务对 PostHog 核心功能非必需、且长迁移在短时间内完成概率极低,这个风险是可接受的。 - 回滚:
attempt_migration_rollback(runner.py)从最后一个已启动操作开始,按逆序遍历执行每个操作的rollback_fn;任一回滚失败则记录错误并停止(防止部分回滚留下不一致状态),全部成功后状态置为RolledBack、进度归零。回滚整体超时受ASYNC_MIGRATIONS_ROLLBACK_TIMEOUT控制。 - 错误处理:
process_error会把错误消息写入AsyncMigrationError记录(关联到迁移外键),记录finished_at,可发送遥测事件与告警邮件(ASYNC_MIGRATIONS_OPT_OUT_EMAILS可关闭),并默认自动触发回滚。以下情况不会自动回滚:显式rollback=False、状态为FailedAtStartup、或开启了ASYNC_MIGRATIONS_DISABLE_AUTO_ROLLBACK。
迁移状态机定义在 models/async_migration.py 的MigrationStatus:NotStarted=0、Running=1、CompletedSuccessfully=2、Errored=3、RolledBack=4、Starting=5(仅 UI 相关)、FailedAtStartup=6。AsyncMigration模型还持久化progress、current_operation_index、current_query_id、celery_task_id、started_at/finished_at、版本区间与parameters(JSON),是断点续跑与 UI 展示的数据基础。
相关配置项一览
异步迁移的行为由环境变量/实例设置控制,默认值见 settings/dynamic_settings.py:
| 配置项 | 默认值 | 作用 |
|---|---|---|
AUTO_START_ASYNC_MIGRATIONS | False | 服务启动时是否自动触发最早未应用的异步迁移 |
ASYNC_MIGRATIONS_DEFAULT_TIMEOUT_SECONDS | 2 * 60 * 60 | SQL 操作默认执行超时(见 settings/async_migrations.py) |
ASYNC_MIGRATIONS_ROLLBACK_TIMEOUT | 30 | 完整回滚的超时时间 |
ASYNC_MIGRATIONS_DISABLE_AUTO_ROLLBACK | False | 是否禁用失败迁移的自动回滚 |
ASYNC_MIGRATIONS_AUTO_CONTINUE | True | Celery worker 崩溃后是否自动断点续跑迁移 |
ASYNC_MIGRATIONS_BLOCK_UPGRADE | True | 存在运行中/出错/必需的迁移时是否阻止升级 |
ASYNC_MIGRATIONS_IGNORE_POSTHOG_VERSION | False | 是否忽略迁移的 PostHog 版本限制(高级) |
ASYNC_MIGRATIONS_OPT_OUT_EMAILS | False | 是否退订迁移完成/失败的邮件通知 |
代码库结构
文档最后给出了各模块的职责划分,对照当前仓库实际路径如下:
| 文档中的模块 | 实际路径 | 职责 |
|---|---|---|
posthog/models/async_migration.py | posthog/models/async_migration.py | Django ORM(Postgres)模型,存储迁移元数据、状态机与查询辅助函数 |
posthog/api/async_migrations.py | posthog/api/async_migration.py | REST API,提供迁移数据查询与启动/停止/回滚触发(测试见 test_async_migrations.py) |
posthog/tasks/async_migrations.py | posthog/tasks/async_migrations.py | Celery 任务:run_async_migration与check_async_migration_health |
posthog/async_migrations/definition.py | posthog/async_migrations/definition.py | 编写迁移所需的基类与操作类型 |
posthog/async_migrations/setup.py | posthog/async_migrations/setup.py | 服务启动时的初始化脚手架与依赖图构建 |
posthog/async_migrations/runner.py | posthog/async_migrations/runner.py | 迁移执行核心:顺序执行操作、回滚、前置检查 |
posthog/async_migrations/utils.py | posthog/async_migrations/utils.py | 不依赖迁移定义的工具函数:执行 SQL、错误处理、强制停止等 |
| — | posthog/async_migrations/status.py | 迁移状态/常量相关辅助 |
| — | posthog/version_requirement.py | ServiceVersionRequirement服务版本约束实现 |
对编写者而言最值得参考的还有两处:真实迁移实现posthog/async_migrations/migrations/(如0005_person_replacing_by_version、0009_minmax_indexes_for_materialized_columns等),以及测试套件posthog/async_migrations/test/(含 test_runner.py、test_definition.py、test_utils.py 与针对具体迁移的测试),测试中大量使用_cases与 mock 来验证执行与回滚路径,是理解框架行为的最佳入口。
小结
编写 PostHog 异步迁移的要点可以归纳为:在posthog/async_migrations/migrations/下按编号命名创建文件,导出继承AsyncMigrationDefinition的Migration类;用operations声明有序的 SQL/通用操作并为关键操作提供回滚;用is_required保证全新实例跳过、healthcheck保证执行环境安全、progress提供可视化进度;在"并发数为 1、版本兼容、依赖完成、健康检查通过"的前提下由 Celery 后台执行,并由 30 分钟一次的健康检查任务兜底 worker 崩溃与执行环境恶化。理解状态机与回滚语义后,你就能安全地为自托管用户编排复杂的数据迁移。
【免费下载链接】posthog:hedgehog: PostHog is the leading platform for building self-driving products. Our developer tools – AI observability, analytics, session replay, flags, experiments, error tracking, logs, and more – capture all the context agents need to diagnose problems, uncover opportunities, and ship fixes. Steer it all from Slack, web, desktop, or the MCP.项目地址: https://gitcode.com/GitHub_Trending/po/posthog
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考