ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

PostHog 悬空 Person 修复指南:基于 Temporal 的 Sync Person Distinct IDs 工作流深入解析

PostHog 悬空 Person 修复指南:基于 Temporal 的 Sync Person Distinct IDs 工作流深入解析 PostHog 悬空 Person 修复指南基于 Temporal 的 Sync Person Distinct IDs 工作流深入解析【免费下载链接】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 的架构中用户身份数据被拆分存储在两套数据库ClickHouse 承载分析查询所需的人员表person、person_distinct_id2PostgreSQL 中的 persons 数据库则承载权威的人员关系。由于写入路径存在竞态条件或历史数据迁移不完整ClickHouse 的person表可能出现悬空记录——人员存在且is_deleted 0却没有对应的person_distinct_id2记录。本篇文章围绕 PostHog 仓库中 Sync Person Distinct IDs 工作流文档 展开结合 workflow.py、activities.py 等源码与测试完整讲解该 Temporal 工作流的孤儿人员分类、输入输出参数、四个 Activity 的实现原理、行为矩阵、CLI 操作以及本地测试方法。读完本文你将掌握如何安全地排查与修复 ClickHouse 中的悬空 Person并理解 PostHog 跨库身份数据同步的工程实践。问题背景ClickHouse 中的悬空 PersonPostHog 将人员身份数据分布在两处ClickHouseperson表与person_distinct_id2表服务于分析查询数据经过 Kafka 异步写入PostgreSQLpersons 数据库posthog_person与posthog_persondistinctid表是身份关系的权威来源。当一次create_person写入与create_person_distinct_id写入之间的时序被打破例如 Kafka 消息乱序、消费者失败、或历史迁移脚本只写了 person 没写 distinct IDClickHouse 中就可能出现SELECT count() FROM person FINAL WHERE team_id %(team_id)s AND is_deleted 0 AND id NOT IN ( SELECT DISTINCT person_id FROM person_distinct_id2 FINAL WHERE team_id %(team_id)s );这类有 Person 但无 Distinct ID的记录就是悬空 Personorphaned person会导致基于 distinct ID 的查询、合并、去重出现偏差。孤儿类别Orphan Categories对照两库数据可以将悬空 Person 划分为三类工作流的处置策略各不相同类别Person 在 ClickHousePerson 在 PostgreSQLDID 在 ClickHouseDID 在 PostgreSQL处置动作可修复Fixable有有无有将 DID 同步回 ClickHouse真正孤儿Truly orphaned有有无无仅报告仅 CH 孤儿CH-only orphans有无无N/A在 ClickHouse 中标记删除可选其中可修复是最常见的场景DID 数据仍然安全地存在于 PostgreSQL只是没有镜像到 ClickHouse真正孤儿则表明连权威库都丢失了身份映射只能报告留存仅 CH 孤儿意味着 PostgreSQL 中根本没有这个人员它很可能是脏数据残留。解决方案工作流总体流程SyncPersonDistinctIdsWorkflowTemporal 注册名为sync-person-distinct-ids按照以下步骤执行在 ClickHouse 中查找所有没有 distinct ID 的孤儿 Person单次查询结果集只有 UUID 元数据体量小在 PostgreSQL 中查询这些 Person UUID 对应的 distinct ID通过 Kafka 将缺失的person_distinct_id2记录同步到 ClickHouse可选将仅存在于 ClickHouse 的孤儿PG 中无数据标记为已删除。从源码看工作流继承了PostHogWorkflow见 start_temporal_workflow.py 中is_named的类型约束并且通过init.py 将工作流与四个 Activity 注册进WORKFLOWS/ACTIVITIES列表供 Temporal Worker 加载。工作流输入参数详解工作流入口为SyncPersonDistinctIdsWorkflowInputs数据类对应源码 workflow.pydataclasses.dataclass class SyncPersonDistinctIdsWorkflowInputs: team_id: int batch_size: int 100 dry_run: bool True # Safe by default delete_ch_only_orphans: bool False # If True, mark CH-only orphans as deleted (requires categorize_orphansTrue) categorize_orphans: bool False # If True, run extra query to distinguish truly orphaned vs CH-only limit: int | None None # Max persons to process (for testing) person_ids: list[str] | None None # Specific person UUIDs to process (for testing) pg_statement_timeout_seconds: int 5各参数说明参数默认值作用team_id必填要处理的团队 ID工作流一次只处理一个团队batch_size100每个批次处理的 Person 数量平衡吞吐与内存dry_runTrue干跑模式只记录日志不做任何写操作默认开启保证安全delete_ch_only_orphansFalse是否将 CH-only 孤儿标记为删除破坏性操作必须显式开启categorize_orphansFalse是否执行额外查询区分真正孤儿与CH-only 孤儿limitNone最多处理多少 Person用于测试person_idsNone指定要处理的 Person UUID 列表用于测试pg_statement_timeout_seconds5PostgreSQL 查询语句超时时间秒源码中有一个容易被忽略的关键校验workflow.py__post_init__强制要求delete_ch_only_orphansTrue时categorize_orphans必须为True否则抛出ValueError。原因正如文档所注只有先分类才能确认待删除的对象确实是PG 中完全不存在的 CH-only 孤儿而不是在 PG 中存在但没有 DID的真正孤儿——后者标记删除将造成不可逆的数据丢失。此外properties_to_log属性会在日志中脱敏输出关键入参person_ids只记录数量而非具体 UUID便于排查运行时问题。工作流结果字段详解工作流返回SyncPersonDistinctIdsWorkflowResultdataclasses.dataclass class SyncPersonDistinctIdsWorkflowResult: team_id: int total_orphaned_persons: int persons_with_pg_distinct_ids: int # Fixable - have DIDs in PG distinct_ids_synced: int persons_without_pg_data: int # Truly orphaned CH-only combined persons_truly_orphaned: int # In PG but no DIDs (only if categorize_orphansTrue) persons_ch_only: int # Not in PG at all (only if categorize_orphansTrue) persons_marked_deleted: int # Only if delete_ch_only_orphansTrue and dry_runFalse dry_run: bool字段含义对应关系total_orphaned_personsClickHouse 中发现的孤儿 Person 总数persons_with_pg_distinct_ids在 PG 中有 DID 的人数可修复distinct_ids_synced实际同步到 ClickHouse 的 distinct ID 条数注意干跑模式下该值也会按将同步的数量返回见下文 Activity 源码persons_without_pg_dataPG 中无 DID 的人数真正孤儿 CH-only 合并persons_truly_orphaned/persons_ch_only仅当categorize_orphansTrue时才有意义的分项统计persons_marked_deleted仅当delete_ch_only_orphansTrue且dry_runFalse时才会大于 0。注意categorize_orphansTrue会触发一次额外的 PostgreSQL 查询将PG 中存在但无 DID与PG 中完全不存在区分开这对报告很有价值但会带来查询开销因此默认关闭。四个 Activity 的源码级剖析工作流的核心逻辑全部封装在四个 Activity 中定义于 activities.py。每个 Activity 都用Heartbeater包裹为长耗时操作提供心跳与进度明细heartbeater.details配合 Temporal 的heartbeat_timeout防止任务被误判为超时。1. find_orphaned_persons单次查询发现孤儿async def find_orphaned_persons(inputs: FindOrphanedPersonsInputs) - FindOrphanedPersonsResult:数据库ClickHouse。职责找出团队内所有没有 distinct ID 的 Person。源码activities.py中有两个值得关注的实现细节其一用argMax/maxGROUP BY取代FINAL去重。注释明确说明FINAL的去重结果取决于合并状态可能不稳定导致有合法 distinct_id 的 Person 被误判为孤儿。因此查询改写为SELECT id AS person_id, any(team_id) AS tid, toString(argMax(created_at, version)) AS created_at, max(version) AS ver FROM person WHERE team_id %(team_id)s [AND id IN %(person_ids)s] GROUP BY id HAVING argMax(is_deleted, version) 0 AND id NOT IN ( SELECT person_id FROM person_distinct_id2 WHERE team_id %(team_id)s GROUP BY person_id, distinct_id HAVING argMax(is_deleted, version) 0 ) ORDER BY created_at ASC [LIMIT %(limit)s] FORMAT JSONEachRow同时子查询的别名tid、ver刻意与列名不同避免 ClickHouse 在WHERE/HAVING/argMax子句中混淆输出别名与列引用。子查询同样按person_id, distinct_id分组并过滤已删除的 DID确保只有当前有效、且没有任何有效 DID 关联的 Person 才被判定为孤儿。其二返回的version是为后续删除做准备。OrphanedPerson数据结构包含person_id、team_id、created_at、version。工作流会维护person_versions字典因为标记删除时 ClickHouse 需要current_version 1才能覆盖旧记录。活动支持两个可选入参limit拼接LIMIT子句与person_ids拼接AND id IN %(person_ids)s后者正是 CLI 中只处理指定 UUID能力的底层实现。2. lookup_pg_distinct_idsPostgreSQL 权威查询与分类async def lookup_pg_distinct_ids(inputs: LookupPgDistinctIdsInputs) - LookupPgDistinctIdsResult:数据库PostgreSQLpersons 数据库。职责为一批 Person UUID 查询其 distinct ID 及其版本号。核心查询activities.pySELECT p.uuid::text as person_uuid, pdi.distinct_id, COALESCE(pdi.version, 0) as version FROM posthog_person p JOIN posthog_persondistinctid pdi ON pdi.person_id p.id AND pdi.team_id p.team_id WHERE p.team_id %s AND p.uuid IN (SELECT unnest(%s::uuid[])) AND p.is_deleted false AND pdi.is_deleted false ORDER BY p.uuid, pdi.id值得注意的实现要点连接条件同时约束team_id并过滤两侧is_deleted false保证只同步当前有效的身份映射通过persons_db_aconnection(writerFalse)走只读连接persons 表位于独立数据库对应 README 中的PERSONS_DB_READER_URL每条语句前执行SET statement_timeout {pg_statement_timeout_seconds}s避免慢查询长期占用数据库连接源码使用psycopg.sql安全拼接字面量结果按person_uuid聚合成PersonDistinctIdMapping其中distinct_id_versions是distinct_id - version的字典——每个 distinct ID 拥有自己的版本号而非取 person 的最大版本这是后续合并正确性的关键测试test_syncs_distinct_ids_with_individual_versions专门验证了这一点。当categorize_orphansTrue且存在persons_not_found时活动会执行第二个查询检查这些 UUID 是否仍存在于posthog_person表SELECT uuid::text FROM posthog_person WHERE team_id %s AND uuid IN (SELECT unnest(%s::uuid[])) AND is_deleted false命中者归入persons_truly_orphanedPG 有 Person 但无 DID未命中者归入persons_ch_onlyPG 完全无此人。LookupPgDistinctIdsResult同时返回mappings、persons_not_found两者合并计数、persons_truly_orphaned、persons_ch_only工作流据此累加各类统计。3. sync_distinct_ids_to_ch经 Kafka 回写 ClickHouseasync def sync_distinct_ids_to_ch(inputs: SyncDistinctIdsToChInputs) - SyncDistinctIdsToChResult:数据库Kafka → ClickHouse。职责把从 PG 取回的 DID 映射写入 ClickHouse 的person_distinct_id2表。干跑模式下dry_runTrue活动逐条记录DRY RUN: Would sync日志包含person_uuid与distinct_id_versions并返回与真实同步相同数量的计数——这正是 README 行为矩阵中dry_runtrue 时 Sync DIDs 显示 No (log)的实现方式也解释了为何工作流结果在干跑时distinct_ids_synced也会等于可修复数量。真实同步时活动遍历每个 mapping 的每个 distinct ID调用 create_person_distinct_idcreate_person_distinct_id( team_idinputs.team_id, distinct_iddistinct_id, person_idmapping.person_uuid, versionversion, is_deletedFalse, )该工具函数通过ClickhouseProducer向KAFKA_PERSON_DISTINCT_IDtopic 生产消息再由 Kafka 消费者落库到person_distinct_id2——这就是经 Kafka 同步的实际链路。由于person_distinct_id2使用ReplacingMergeTree且以version去重只要保留 PG 中的原始版本号重复执行工作流不会产生重复或错误覆盖从而实现幂等。4. mark_ch_only_orphans_deleted版本感知的删除标记async def mark_ch_only_orphans_deleted(inputs: MarkChOnlyOrphansDeletedInputs) - MarkChOnlyOrphansDeletedResult:数据库Kafka → ClickHouse。职责将确认是 CH-only 的孤儿 Person 标记为已删除。干跑模式同样只记录DRY RUN: Would mark as deleted日志。真实执行时对每个 person 调用 create_person并以current_version 1作为新版本、is_deletedTrue写入create_person( uuidperson_uuid, team_idinputs.team_id, versioncurrent_version 1, is_deletedTrue, )版本必须 1 的原因person表同样是ReplacingMergeTree以version决定保留哪条记录。若写入的删除版本不大于当前版本旧的非删除记录会胜出删除将被忽略。测试 test_workflow.py 专门构造了版本为1、1234、123456789的三种 Person断言删除后每条记录的版本恰为原版本 1且is_deleted 1。工作流编排一次发现 批量处理工作流主循环workflow.py的设计核心是避免低效的 OFFSET 分页先执行单次find_orphaned_persons查询一次性取回全部孤儿 UUID结果集只是 UUID 元数据体量小在for i in range(0, len(person_uuids), inputs.batch_size)循环中按batch_size分批执行 PG 查询与 CH 写入——这两步才是真正昂贵的部分每批内先lookup_pg_distinct_ids若存在mappings则执行sync_distinct_ids_to_ch若delete_ch_only_orphans且本批存在persons_ch_only则构造delete_versions缺失时回退版本 0并执行mark_ch_only_orphans_deleted。对比 OFFSET 分页N 次查询、每次扫描 O(n) 行的方式本方案是1 次发现查询 分批处理整体扫描成本更低。所有execute_activity调用都配置了统一的RetryPolicyinitial_interval10s、maximum_interval2min、maximum_attempts5以及各自的start_to_close_timeout/heartbeat_timeout如 PG 查询 5 分钟/30 秒心跳CH 写入 10 分钟/60 秒心跳。行为矩阵四个开关的组合语义dry_run与delete_ch_only_orphans两个开关决定工作流的实际行为dry_rundelete_ch_only_orphans同步 DID标记孤儿删除truefalse否仅日志否truetrue否仅日志否仅日志falsefalse是否falsetrue是是再次强调delete_ch_only_orphanstrue必须搭配categorize_orphanstrue源码__post_init__强制校验确保只删除确认在 PG 中完全不存在的 CH-only 孤儿而绝不误伤在 PG 中存在但无 DID的真正孤儿。CLI 使用指南工作流通过 start_temporal_workflow 管理命令启动格式为python manage.py start_temporal_workflow WORKFLOW JSON 输入。命令会从注册的WORKFLOWS列表中按名称匹配本工作流注册名为sync-person-distinct-ids解析 JSON 输入并以--workflow-id默认随机 UUID、--task-queue、--max-attempts等参数启动 Temporal 工作流。干跑默认安全# 干跑默认——仅报告将要执行的操作 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 2} # 带 limit 的干跑——只处理前 10 个孤儿测试用 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 2, limit: 10} # 针对指定 Person 的干跑测试用 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 2, person_ids: [uuid-1, uuid-2]} # 带删除标记预览的干跑——同时展示哪些 CH-only 孤儿会被标记删除 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 2, delete_ch_only_orphans: true, categorize_orphans: true} # 带分类的干跑——分别报告真正孤儿与 CH-only 孤儿 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 2, categorize_orphans: true}生产执行# 生产仅同步不标记 CH-only 孤儿删除 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 2, dry_run: false} \ --workflow-id sync-person-distinct-ids-team-2 # 生产同步并将 CH-only 孤儿标记为删除 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 2, dry_run: false, delete_ch_only_orphans: true, categorize_orphans: true} \ --workflow-id sync-person-distinct-ids-team-2--workflow-id的作用值得注意默认随机 UUID 意味着每次都启动新工作流指定固定 ID 后如sync-person-distinct-ids-team-2若该 ID 已有运行中或已完成的工作流则不会重复启动WorkflowIDReusePolicy.ALLOW_DUPLICATE_FAILED_ONLY从而限制同一团队并发执行——这与 README 设计决策中单团队单工作流实例的隔离思路一致。前置条件是本地必须已运行 Temporal Server 并启动了相应 Workerpython manage.py start_temporal_worker任务队列需配置为运行本工作流的队列。本地测试构造孤儿数据并验证仓库提供了专门的测试数据管理命令 setup_orphan_test_data无需手工编写 SQL 即可构造三类孤儿。构造测试数据# 创建默认测试孤儿3 个可修复、2 个真正孤儿、2 个 CH-only python manage.py setup_orphan_test_data --team-id 1 # 自定义数量 python manage.py setup_orphan_test_data --team-id 1 \ --fixable 5 \ --truly-orphaned 3 \ --ch-only 2 # 自定义 distinct ID 前缀 python manage.py setup_orphan_test_data --team-id 1 --prefix my-test命令实现细节对应源码可修复孤儿在 PG 的posthog_personposthog_persondistinctid中写入 person 与 DID通过insert_seed_person/insert_seed_distinct_id但只在 ClickHouseperson表插入 Person、刻意不插入 DID真正孤儿只在 PG 建 Person 不建 DIDCH-only 孤儿则只在 ClickHouse 插入 Person、完全不在 PG 建记录。所有 PG 写入走persons_db_connection(writerTrue, autocommitTrue)直接落库避免 ORM 辅助方法同步镜像到 ClickHouse 破坏孤儿状态。运行工作流# 干跑——查看将被同步/删除的内容 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 1} # 实际同步可修复孤儿 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 1, dry_run: false} # 同步 将 CH-only 孤儿标记为删除 python manage.py start_temporal_workflow sync-person-distinct-ids \ {team_id: 1, dry_run: false, delete_ch_only_orphans: true, categorize_orphans: true}验证结果-- 检查 ClickHouse 中剩余的孤儿 SELECT id, team_id, is_deleted, version FROM person FINAL WHERE team_id 1 AND is_deleted 0 AND id NOT IN ( SELECT DISTINCT person_id FROM person_distinct_id2 FINAL WHERE team_id 1 ); -- 检查已同步的 distinct ID SELECT person_id, distinct_id, version, is_deleted FROM person_distinct_id2 FINAL WHERE team_id 1 AND distinct_id LIKE test-orphan%;清理测试数据python manage.py setup_orphan_test_data --team-id 1 --cleanup清理逻辑会从 PG 中删除带test_type标记的测试 Person 与其 DID--prefix需匹配创建时使用的前缀源码对%/_等 LIKE 通配符做了转义处理并在 ClickHouse 中以version 100将测试 Person 标记删除。注意CH-only 孤儿在 PG 中无记录、无法通过test_type定位命令会提示改用delete_ch_only_orphanstrue的运行方式来清理。测试场景速查表场景命令JSON 输入预期结果干跑{team_id: 1}记录数量日志无任何变更分类{team_id: 1, categorize_orphans: true}分别报告真正孤儿与 CH-only 孤儿仅同步{team_id: 1, dry_run: false}可修复孤儿获得 DID 同步同步 删除{team_id: 1, dry_run: false, delete_ch_only_orphans: true, categorize_orphans: true}DID 同步 CH-only 标记删除限制数量{team_id: 1, limit: 2}只处理前 2 个孤儿指定 Person{team_id: 1, person_ids: [uuid-1]}只处理指定 UUID测试覆盖从 Activity 到端到端仓库在 posthog/temporal/tests/sync_person_distinct_ids/ 下提供了两层测试均标记pytest.mark.persons_db_direct说明它们直接读写真实的 persons 数据库与 ClickHousetest_activities.py逐项验证每个 Activity 的行为——find_orphaned_persons能发现孤儿、跳过有 DID 的 Person、跳过已删除的 Person、遵守limitlookup_pg_distinct_ids能返回多 DID 及各自的独立版本、正确区分真正孤儿与 CH-onlysync_distinct_ids_to_ch干跑不落库、真实写入保留每个 DID 的原始版本mark_ch_only_orphans_deleted干跑不删除、对高版本 Person 也能以版本 1成功标记删除。文件末尾的TestEndToEndOrphanCategories一次性构造三类孤儿并断言分类结果互不重叠。test_workflow.py使用WorkflowEnvironment.start_time_skipping()与真实 Worker 执行完整工作流构造 100 个孤儿40 可修复 30 真正孤儿 30 CH-only以BATCH_SIZE30强制触发多批次验证干跑不产生变更、同步模式后仅剩真正孤儿、同步 删除模式后仅剩真正孤儿以及版本感知删除的正确性版本 1/1234/123456789 分别变为 1。这些测试同时印证了 README 中的设计决策每个 distinct ID 保留独立版本distinct_id_versions字典、删除必须版本递增、干跑与分类默认关闭以保安全。设计决策总结README 列出的 10 条设计决策与源码一一对应是理解本工作流工程取舍的关键单团队单工作流每次运行只处理一个team_id不同团队分开运行实现简单且隔离性好CLI 通过--workflow-id进一步限制并发批大小 100平衡吞吐与内存测试中以 30 验证多批处理路径幂等性person_distinct_id2使用ReplacingMergeTree version 去重重复执行安全默认干跑dry_runTrue只记录日志不做变更删除标记为可选项破坏性操作必须显式传delete_ch_only_orphansTrue分类为可选项额外查询区分真正孤儿与 CH-only默认关闭以避免开销心跳机制所有 Activity 使用Heartbeater支撑长操作独立 persons 数据库PG 查询通过persons_db_aconnection连接PERSONS_DB_READER_URL对应的独立 persons 数据库按 DID 维护版本每个 distinct ID 独立版本同步时原样保留防止后续合并因版本过低被忽略版本感知删除标记删除时使用current_version 1确保 ClickHouse 的 ReplacingMergeTree 采纳删除记录。结语Sync Person Distinct IDs 工作流是 PostHog 应对跨库身份数据漂移的典型工程实践以 ClickHouse 单次发现查询定位问题集合以 PostgreSQL 为权威源补齐数据以 Kafka 异步通道完成回写并以默认干跑 显式破坏 版本感知三重机制保证操作安全。无论是排查历史迁移遗留问题还是理解 PostHog 人员身份数据的写入链路本文涉及的源码workflow.py、activities.py、util.py与测试test_workflow.py、test_activities.py都是可直接查阅的权威参考。【免费下载链接】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),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进