PostHog 数据建模 Temporal 工作流解析:物化视图与 DAG 编排的工程实践
2026/9/13 12:15:12 网站建设 项目流程

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,共三个工作流类型:MaterializeViewWorkflowExecuteDAGWorkflowEnrichViewSemanticsWorkflow(后者负责视图语义描述的异步刷新)。活动则统一从 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_EXISTINGWorkflowIDReusePolicy.ALLOW_DUPLICATE,保证同一节点的并发触发会复用已有运行;
  • 工作流级的RetryPolicy(maximum_attempts=1),把重试策略完全交给工作流内部的活动去处理;
  • 任务队列为settings.DATA_MODELING_TASK_QUEUE

单视图物化:MaterializeViewWorkflow

posthog/temporal/data_modeling/workflows/materialize_view.py 定义了名为data-modeling-materialize-viewMaterializeViewWorkflow。它负责单个视图/物化视图的完整物化生命周期,既可以由用户点击"立即物化"直接触发,也可以作为 DAG 编排工作流的子工作流被调用。

输入与输出协议

工作流的输入MaterializeViewWorkflowInputs是冻结的 dataclass,字段如下:

字段类型默认值说明
team_idint必填拥有该节点的团队 ID
dag_idstr必填节点所属 DAG
node_idstr必填待物化节点的 UUID
managed_warehouse_onlyboolFalse仅走托管数仓(Managed Warehouse)引擎
dangerously_execute_raw_sqlboolFalse允许执行原始 SQL(需显式开启)
manually_triggered_by_idint | NoneNone发起本次运行的用户(人工触发时)
duckgres_onlyboolFalse遗留字段:旧负载包含该字段,删除会导致部署后重放失败,因此保留

输出MaterializeViewWorkflowResult包含:job_id(本次运行创建的DataModelingJob记录 ID)、node_idrows_materialized(写入 Delta 表的行数)、duration_secondsquality_blocking_failures(阻塞发布的错误级检查失败数,None 表示未执行门禁审计)、quality_audited(本次运行是否已被检查套件覆盖,供 DAG 收尾扫描避免重复执行)。

执行流程

工作流的主流程可以概括为五步:

  1. 创建任务记录:通过create_data_modeling_job_activity写入DataModelingJob行,记录team_idnode_iddag_id、父工作流 ID 与手动触发者;
  2. 执行 HogQL 查询并写入 Delta Lake:调用materialize_view_activity运行查询、按批次写出 parquet 文件并生成 Delta 表;
  3. 准备可查询表:通过prepare_queryable_table_activity(或数据质量门禁下的 stage/publish 三连)把文件整理为可查询结构;
  4. 收尾:成功后由succeed_materialization_activity更新节点与任务状态,失败则由fail_materialization_activity标记失败并记录错误;
  5. 触发旁路副作用(成功后):包括视图语义描述刷新(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_timeoutheartbeat_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-dagExecuteDAGWorkflow。它负责编排 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_nodessuccessful_nodesfailed_nodesskipped_nodesduration_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_activitycreate_data_modeling_job.py创建任务记录 / 记录被跳过的节点
get_dag_structure_activityget_dag_structure.py从数据库读取 DAG 的节点、边、可执行节点、临时节点与挂起节点
materialize_view_activitymaterialize_view.py执行 HogQL 查询并写 Delta 表(v2 主路径)
materialize_view_duckgres_activity/materialize_view_managed_warehouse_activitymaterialize_view_managed_warehouse.pyDuckgres / 托管数仓阴影(shadow)物化
prepare_queryable_table_activity/stage_queryable_files_activity/publish_queryable_table_activityprepare_queryable_table.py准备可查询表 / 暂存文件 / 发布
quality_block_materialization_activityquality_block_materialization.py因数据质量检查失败而阻止发布
succeed_materialization_activity/fail_materialization_activitysucceed_materialization.py/fail_materialization.py成功 / 失败收尾
preempt_dag_run_activitypreempt_dag_run.pyDAG 运行开始前清理遗留的脏任务状态
notify_materialization_failure.pynotify_materialization_failure.py失败通知
enrich_view_semantics_activityenrich_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_datetimesarrow_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 结构:可执行节点限定为VIEWMAT_VIEWENDPOINT三类且排除软删除的保存查询;VIEW类型视为临时(ephemeral)节点——无需物化,直接标记成功;挂起节点按引擎维度组织成字典(suspended_nodes[engine])。

退役的 v1 后端:一次谨慎的删除

CLAUDE.md 明确记录了 v1 退役的教训:run_workflow.py与其服务的data-modeling-run每查询调度已删除。删除工作流类型必须放在最后——因为指向已注销类型的调度不会响亮地失败,它会持续触发、工作流任务失败,却不写入任何任务行。因此只有在两个 region 都不再有data-modeling-run调度后,该工作流类型才被注销。

v1 有两件"遗产"被刻意保留:

  1. resolve_log_source仍解析 v1 工作流 ID 的形状:否则所有切换前的运行都会丢失日志;
  2. v1 的DataModelingJob:它们是切换前运行情况的唯一记录。

这一设计体现了 Temporal 工作流演进的一个通用原则:历史事件(history)是不可重写的,删除类型/字段必须考虑在飞运行与既有历史的重放兼容性。当前代码中同类思想随处可见,例如MaterializeViewWorkflowInputs.duckgres_only字段的注释("Old workflow payloads contain this field, so removing it would prevent replay after deployment"),以及QUALITY_AUDIT_PATCHCDP_VIEW_TRIGGER_PATCHMANAGED_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.pymodeling.py),但其 viewset 仍保留在products/data_warehouse/下。

对照实际目录结构可以看到,数据建模的领域模型(NodeEdgeDAGDataModelingJobDataWarehouseSavedQuery等)分布在 products/data_modeling/backend/models/ 下,通过 facade/api.py 这层 facade 向外暴露逻辑能力(如start_node_materialization、增量配置、挂起管理等),而 Temporal 工作流侧只依赖这些 facade 与活动,形成清晰的模块化分层。

小结:一张完整的调用链地图

把整条链路串起来看:

  1. 用户触发节点materialize或保存查询runstart_node_materialization启动data-modeling-materialize-view
  2. 用户触发 DAGrun→ 启动data-modeling-execute-dag
  3. ExecuteDAGWorkflow先跑preempt_dag_run_activity清理脏状态,再经get_dag_structure_activity读取结构,Kahn 拓扑排序分层后按层并行(信号量限流 10)启动MaterializeViewWorkflow子工作流;
  4. MaterializeViewWorkflow创建DataModelingJob→ 执行 HogQL 物化写 Delta 表 → 经数据质量门禁(可选)→ 准备并发布可查询表 → 成功/失败收尾 → 异步触发语义刷新、属性同步与 CDP 触发等旁路副作用;
  5. 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询