PostHog 数据建模 Temporal 工作流解析:物化视图与 DAG 编排的工程实践
【免费下载链接】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
导读
PostHog 的数据建模(data modeling)能力依赖一套运行在 Temporal 之上的工作流系统:它将用户在数据仓库中定义的保存查询(saved query)物化为可查询的 Delta Lake 表,并把存在依赖关系的多个物化任务组织成有向无环图(DAG)按拓扑序执行。本文以 posthog/temporal/data_modeling/CLAUDE.md 为核心骨架,结合源码深入讲解这套工作流的目录布局、两个核心工作流的输入输出与执行细节、v2 后端对 v1 的演进原因,以及团队在跨模块边界上的职责划分。读完本文,你将掌握 PostHog 数据建模物化的完整调用链:从节点运行入口到单视图物化工作流,再到 DAG 编排工作流与活动(Activity)体系。
代码布局:一套专注于物化的 Temporal 工作流树
数据建模 Temporal 工作流位于 posthog/temporal/data_modeling/,其职责一句话可以概括:物化数据建模保存查询。当前树内只有 v2 后端,v1 早已被删除(详见后文"退役的 v1 后端"一节)。
目录布局如下:
workflows/materialize_view.py:物化一个保存查询(单节点物化工作流);workflows/execute_dag.py:运行一个 DAG,或 DAG 中某个 tier 的节点子集(编排工作流);activities/*.py:全部活动(Activity)实现,工作流本身不触碰数据库与外部服务;- 入口约定:Node 的
run启动data-modeling-execute-dag工作流;Node 的materialize与保存查询的run则通过start_node_materialization启动data-modeling-materialize-view工作流。
工作流的注册与导出集中在 posthog/temporal/data_modeling/workflows/init.py,共三个工作流类型:MaterializeViewWorkflow、ExecuteDAGWorkflow、EnrichViewSemanticsWorkflow(后者负责视图语义描述的异步刷新)。活动则统一从 posthog/temporal/data_modeling/activities/init.py 导出,包括创建/失败/成功标记数据建模任务、获取 DAG 结构、执行 HogQL 物化、准备可查询表、数据质量拦截、通知失败等十余个活动。
关键入口:start_node_materialization
CLAUDE.md 提到的start_node_materialization实现在 products/data_modeling/backend/logic/node_materialization.py,它同时服务于 Node 的materialize与保存查询的run两个动作。从源码可以看出几个关键约定:
- 显式触发的运行(
resume=True)会先解除节点挂起(unsuspend_nodes),因为"用户主动再试一次"意味着重新获得一次全新的失败窗口; resume=False用于读取流量触发的自动修复,此时即便节点处于挂起状态,运行照常启动,但不会清除挂起标记——防止请求流量撑开熔断器;- 工作流 ID 固定为
materialize-view-{node.id},配合WorkflowIDConflictPolicy.USE_EXISTING与WorkflowIDReusePolicy.ALLOW_DUPLICATE,保证同一节点的并发触发会复用已有运行; - 工作流级的
RetryPolicy(maximum_attempts=1),把重试策略完全交给工作流内部的活动去处理; - 任务队列为
settings.DATA_MODELING_TASK_QUEUE。
单视图物化:MaterializeViewWorkflow
posthog/temporal/data_modeling/workflows/materialize_view.py 定义了名为data-modeling-materialize-view的MaterializeViewWorkflow。它负责单个视图/物化视图的完整物化生命周期,既可以由用户点击"立即物化"直接触发,也可以作为 DAG 编排工作流的子工作流被调用。
输入与输出协议
工作流的输入MaterializeViewWorkflowInputs是冻结的 dataclass,字段如下:
| 字段 | 类型 | 默认值 | 说明 |
|---|---|---|---|
team_id | int | 必填 | 拥有该节点的团队 ID |
dag_id | str | 必填 | 节点所属 DAG |
node_id | str | 必填 | 待物化节点的 UUID |
managed_warehouse_only | bool | False | 仅走托管数仓(Managed Warehouse)引擎 |
dangerously_execute_raw_sql | bool | False | 允许执行原始 SQL(需显式开启) |
manually_triggered_by_id | int | None | None | 发起本次运行的用户(人工触发时) |
duckgres_only | bool | False | 遗留字段:旧负载包含该字段,删除会导致部署后重放失败,因此保留 |
输出MaterializeViewWorkflowResult包含:job_id(本次运行创建的DataModelingJob记录 ID)、node_id、rows_materialized(写入 Delta 表的行数)、duration_seconds、quality_blocking_failures(阻塞发布的错误级检查失败数,None 表示未执行门禁审计)、quality_audited(本次运行是否已被检查套件覆盖,供 DAG 收尾扫描避免重复执行)。
执行流程
工作流的主流程可以概括为五步:
- 创建任务记录:通过
create_data_modeling_job_activity写入DataModelingJob行,记录team_id、node_id、dag_id、父工作流 ID 与手动触发者; - 执行 HogQL 查询并写入 Delta Lake:调用
materialize_view_activity运行查询、按批次写出 parquet 文件并生成 Delta 表; - 准备可查询表:通过
prepare_queryable_table_activity(或数据质量门禁下的 stage/publish 三连)把文件整理为可查询结构; - 收尾:成功后由
succeed_materialization_activity更新节点与任务状态,失败则由fail_materialization_activity标记失败并记录错误; - 触发旁路副作用(成功后):包括视图语义描述刷新(
EnrichViewSemanticsWorkflow子工作流)、person/account 属性同步子工作流、CDP 生产者工作流等,全部以ParentClosePolicy.ABANDON隔离启动,任何失败都不会反过来拖垮物化本身。
失败处理与重试策略
源码中定义了一组不可重试错误类型NON_RETRYABLE_ERRORS(materialize_view.py),包括:
CHQueryErrorMemoryLimitExceeded CannotCoerceColumnException InvalidNodeTypeException NodeNotFoundException EmptyHogQLResponseColumnsError这些错误代表"查询或数据本身有问题"而非瞬时故障,重试没有意义。活动调用普遍采用RetryPolicy(maximum_attempts=3, initial_interval=10s, maximum_interval=5min)搭配start_to_close_timeout与heartbeat_timeout(物化活动为 20 分钟超时、2 分钟心跳)的配置。此外,工作流对取消做了专门处理:_is_cancellation会识别直接取消与活动/子工作流包装后的取消,一旦判定为取消,错误消息记为 "Workflow was cancelled",且不会把取消当成业务失败上报。
数据质量门禁:stage / audit / publish
QUALITY_AUDIT_PATCH = "data-quality-audit-2026-08"覆盖了数据质量功能在此处新增的三种命令(stage/audit/publish 三连与 warn 模式下的检查套件子工作流)。门禁模式有三种:QUALITY_AUDIT_GATE(阻塞)、QUALITY_AUDIT_WARN(警告)、QUALITY_AUDIT_SKIP(跳过)。在 GATE 模式下,物化结果先经stage_queryable_files_activity暂存,再运行检查套件子工作流data-quality-gate-{job_id},若存在阻塞性失败(staged_verdict非空),则调用quality_block_materialization_activity阻止发布并返回带quality_blocking_failures的结果——DAG 编排层据此把该节点标记为质量失败。值得注意的是:审计管道自身故障(如检查工作流抛错)不会阻塞发布,因为"检查管道坏了"不等于"数据有罪",此时仍会发布并把节点留给 DAG 的收尾扫描再次尝试。
DAG 编排:ExecuteDAGWorkflow
posthog/temporal/data_modeling/workflows/execute_dag.py 定义了名为data-modeling-execute-dag的ExecuteDAGWorkflow。它负责编排 DAG 内所有节点的物化,核心流程为:获取 DAG 结构 → 计算依赖层级 → 逐层并行执行子工作流 → 统计成败并处理跳过。
拓扑排序与层级执行
_dag_execution_levels使用Kahn 拓扑排序把可执行节点划分为多个层级(execute_dag.py):
- 每个节点的入度 = 其上游中仍在本轮执行集合内的节点数(交集过滤掉用户未请求的节点);
- 每轮取出入度为 0 的节点组成一个 level;
- 若某轮取不出任何节点且集合非空,抛出
EmptyDAGOrCycleError,携带各问题节点的入度、依赖与未满足依赖明细,便于定位环。
并发控制与失败传播
- 全 DAG 使用
MAX_CONCURRENT_CHILDREN = 10的信号量限制并发子工作流数量,以"滑动窗口"方式跨层级生效,做托管数仓与 ClickHouse 基础设施的友好邻居; - 每一层级内部通过
asyncio.gather并行启动该层所有节点的MaterializeViewWorkflow子工作流,子工作流 ID 形如materialize-view-{dag_id}-{node_id}-{start_time.isoformat()},ParentClosePolicy.REQUEST_CANCEL保证父级取消时子级一并取消; - 失败传播:某层节点失败后,其全部下游节点会被跳过(
skipped),跳过原因形如Upstream node {blocked_id} failed/failed data quality checks/suspended;被跳过的节点通过record_skipped_data_modeling_jobs_activity记录为跳过任务,跳过的上游名称数量受UPSTREAM_NAMES_IN_SKIP_REASON限制; - 挂起(suspended)节点同样被跳过,跳过原因为 "Node suspended after repeated materialization failures"——这是节点因反复失败被熔断的机制;
- DAG 结束后,若存在失败节点,调用
notify_dag_materialization_failures_activity通知相关方。
结果汇总与度量
ExecuteDAGResult汇总scheduled_nodes、successful_nodes、failed_nodes、skipped_nodes、duration_seconds与逐节点的NodeResult明细。工作流会按四类状态(completed/skipped/failed/partial_failure)上报dag_finished指标,并记录 DAG 时长与成功/失败/跳过节点数指标(定义于 posthog/temporal/data_modeling/metrics.py)。DAG 编排还承担收尾的数据质量检查:对本轮已成功但未在节点级做过审计的节点,通过_run_data_quality_checks以 ABANDON 方式启动检查套件(data-quality-run-suite-{dag_id}-{run_id}),且先用门禁活动询问"是否真的需要检查",避免无检查项的团队为每次物化都付出子工作流成本。
活动体系:Activities 一览
工作流保持确定性,所有外部 I/O 都下沉到 activities/ 下的活动中。从__init__.py的导出清单可以看到完整活动面:
| 活动 | 文件 | 职责 |
|---|---|---|
create_data_modeling_job_activity/record_skipped_data_modeling_jobs_activity | create_data_modeling_job.py | 创建任务记录 / 记录被跳过的节点 |
get_dag_structure_activity | get_dag_structure.py | 从数据库读取 DAG 的节点、边、可执行节点、临时节点与挂起节点 |
materialize_view_activity | materialize_view.py | 执行 HogQL 查询并写 Delta 表(v2 主路径) |
materialize_view_duckgres_activity/materialize_view_managed_warehouse_activity | materialize_view_managed_warehouse.py | Duckgres / 托管数仓阴影(shadow)物化 |
prepare_queryable_table_activity/stage_queryable_files_activity/publish_queryable_table_activity | prepare_queryable_table.py | 准备可查询表 / 暂存文件 / 发布 |
quality_block_materialization_activity | quality_block_materialization.py | 因数据质量检查失败而阻止发布 |
succeed_materialization_activity/fail_materialization_activity | succeed_materialization.py/fail_materialization.py | 成功 / 失败收尾 |
preempt_dag_run_activity | preempt_dag_run.py | DAG 运行开始前清理遗留的脏任务状态 |
notify_materialization_failure.py | notify_materialization_failure.py | 失败通知 |
enrich_view_semantics_activity | enrich_view_semantics.py | 视图语义描述刷新 |
materialize_view_activity的实现要点
activities/materialize_view.py 是物化的核心实现,值得关注的工程细节包括:
- 并发限制:模块级
asyncio.Semaphore(MAX_CONCURRENT_CLICKHOUSE_QUERIES)(值为 10)限制单 worker 上所有活动共享的 ClickHouse 并发查询数; - 类型转换:ClickHouse 的
DateTime/DateTime64/Date/UUID/ENUM/IPv4/IPv6/JSON等类型在写入 Arrow 批次前会经_transform_date_and_datetimes与arrow_type_conversion映射转换为 Arrow/Delta 可表达的类型;高精度 decimal 列降级为decimal128(38, 37); - schema 一致性:
_force_nullable把每列都标记为可空,确保跨批次 schema 一致,避免 delta-rs 的 DataFusion 写入器因大小写敏感而破坏personId这类驼峰列名; - 增量物化:受
data-modeling-incremental-views特性开关控制,增量路径通过_resolve_write_plan判定(首次运行/定义变更/无可用水位线时退回全量重建),并按水位线窗口注入过滤条件、按唯一键 upsert;写计划的原因会随任务暴露,让意外昂贵的运行"自解释"; - 零行结果:查询返回零行时写出仅含 schema 的空 parquet(
_write_empty_parquet_for_zero_rows),保证空表也可查询、物化不会留下无表可查的模型; - CDP 行暂存:
_CDPRowSink以尽力而为方式把写入的行暂存给 CDP 订阅者,暂存失败时丢弃整轮暂存(部分暂存比没有更糟),但绝不失败物化本身。
DAG 结构活动
activities/get_dag_structure.py 从数据库读取 DAG 结构:可执行节点限定为VIEW、MAT_VIEW、ENDPOINT三类且排除软删除的保存查询;VIEW类型视为临时(ephemeral)节点——无需物化,直接标记成功;挂起节点按引擎维度组织成字典(suspended_nodes[engine])。
退役的 v1 后端:一次谨慎的删除
CLAUDE.md 明确记录了 v1 退役的教训:run_workflow.py与其服务的data-modeling-run每查询调度已删除。删除工作流类型必须放在最后——因为指向已注销类型的调度不会响亮地失败,它会持续触发、工作流任务失败,却不写入任何任务行。因此只有在两个 region 都不再有data-modeling-run调度后,该工作流类型才被注销。
v1 有两件"遗产"被刻意保留:
resolve_log_source仍解析 v1 工作流 ID 的形状:否则所有切换前的运行都会丢失日志;- v1 的
DataModelingJob行:它们是切换前运行情况的唯一记录。
这一设计体现了 Temporal 工作流演进的一个通用原则:历史事件(history)是不可重写的,删除类型/字段必须考虑在飞运行与既有历史的重放兼容性。当前代码中同类思想随处可见,例如MaterializeViewWorkflowInputs.duckgres_only字段的注释("Old workflow payloads contain this field, so removing it would prevent replay after deployment"),以及QUALITY_AUDIT_PATCH、CDP_VIEW_TRIGGER_PATCH、MANAGED_WAREHOUSE_NAMING_PATCH等一整套workflow.patched()演进标记(materialize_view.py)——用 Temporal 的版本标记(patch)为滚动部署期间的新旧历史分支保留兼容路径。
职责边界:Scoping
CLAUDE.md 记录了本模块与其他团队的代码边界,这也是贡献者(无论人类还是 Agent)修改代码时必须遵守的约束:
products/data_warehouse/由另一团队拥有,但其中的保存查询表面(presentation/views/saved_query.py)归数据建模团队变更;该目录树下的其他任何改动都需要对方团队评审;DataModelingJob模型位于 products/data_modeling/backend/models/(具体见data_modeling_job.py与modeling.py),但其 viewset 仍保留在products/data_warehouse/下。
对照实际目录结构可以看到,数据建模的领域模型(Node、Edge、DAG、DataModelingJob、DataWarehouseSavedQuery等)分布在 products/data_modeling/backend/models/ 下,通过 facade/api.py 这层 facade 向外暴露逻辑能力(如start_node_materialization、增量配置、挂起管理等),而 Temporal 工作流侧只依赖这些 facade 与活动,形成清晰的模块化分层。
小结:一张完整的调用链地图
把整条链路串起来看:
- 用户触发节点
materialize或保存查询run→start_node_materialization启动data-modeling-materialize-view; - 用户触发 DAG
run→ 启动data-modeling-execute-dag; ExecuteDAGWorkflow先跑preempt_dag_run_activity清理脏状态,再经get_dag_structure_activity读取结构,Kahn 拓扑排序分层后按层并行(信号量限流 10)启动MaterializeViewWorkflow子工作流;MaterializeViewWorkflow创建DataModelingJob→ 执行 HogQL 物化写 Delta 表 → 经数据质量门禁(可选)→ 准备并发布可查询表 → 成功/失败收尾 → 异步触发语义刷新、属性同步与 CDP 触发等旁路副作用;- DAG 汇总节点结果、上报指标、启动收尾质量检查并在有失败时通知相关人员。
这套体系以"工作流确定性 + 活动可重试 + 历史重放兼容 + 旁路副作用隔离"为设计支柱,是 PostHog 数据仓库产品中数据建模能力的运行底座。想要深入了解实现细节的读者,可以继续阅读 materialize_view.py、execute_dag.py 以及 activities/ 下的各活动文件。
【免费下载链接】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),仅供参考