
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_materializationCLAUDE.md 提到的start_node_materialization实现在 products/data_modeling/backend/logic/node_materialization.py它同时服务于 Node 的materialize与保存查询的run两个动作。从源码可以看出几个关键约定显式触发的运行resumeTrue会先解除节点挂起unsuspend_nodes因为用户主动再试一次意味着重新获得一次全新的失败窗口resumeFalse用于读取流量触发的自动修复此时即便节点处于挂起状态运行照常启动但不会清除挂起标记——防止请求流量撑开熔断器工作流 ID 固定为materialize-view-{node.id}配合WorkflowIDConflictPolicy.USE_EXISTING与WorkflowIDReusePolicy.ALLOW_DUPLICATE保证同一节点的并发触发会复用已有运行工作流级的RetryPolicy(maximum_attempts1)把重试策略完全交给工作流内部的活动去处理任务队列为settings.DATA_MODELING_TASK_QUEUE。单视图物化MaterializeViewWorkflowposthog/temporal/data_modeling/workflows/materialize_view.py 定义了名为data-modeling-materialize-view的MaterializeViewWorkflow。它负责单个视图/物化视图的完整物化生命周期既可以由用户点击立即物化直接触发也可以作为 DAG 编排工作流的子工作流被调用。输入与输出协议工作流的输入MaterializeViewWorkflowInputs是冻结的 dataclass字段如下字段类型默认值说明team_idint必填拥有该节点的团队 IDdag_idstr必填节点所属 DAGnode_idstr必填待物化节点的 UUIDmanaged_warehouse_onlyboolFalse仅走托管数仓Managed Warehouse引擎dangerously_execute_raw_sqlboolFalse允许执行原始 SQL需显式开启manually_triggered_by_idint | NoneNone发起本次运行的用户人工触发时duckgres_onlyboolFalse遗留字段旧负载包含该字段删除会导致部署后重放失败因此保留输出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_ERRORSmaterialize_view.py包括CHQueryErrorMemoryLimitExceeded CannotCoerceColumnException InvalidNodeTypeException NodeNotFoundException EmptyHogQLResponseColumnsError这些错误代表查询或数据本身有问题而非瞬时故障重试没有意义。活动调用普遍采用RetryPolicy(maximum_attempts3, initial_interval10s, maximum_interval5min)搭配start_to_close_timeout与heartbeat_timeout物化活动为 20 分钟超时、2 分钟心跳的配置。此外工作流对取消做了专门处理_is_cancellation会识别直接取消与活动/子工作流包装后的取消一旦判定为取消错误消息记为 Workflow was cancelled且不会把取消当成业务失败上报。数据质量门禁stage / audit / publishQUALITY_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 编排ExecuteDAGWorkflowposthog/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_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_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为滚动部署期间的新旧历史分支保留兼容路径。职责边界ScopingCLAUDE.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用户触发 DAGrun→ 启动data-modeling-execute-dagExecuteDAGWorkflow先跑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),仅供参考