拓冰建站拓冰建站
首页 / 资讯中心 / 正文

PostHog 信号发射管道(Signal Emission Pipeline)架构与接入实战

PostHog 信号发射管道Signal Emission Pipeline架构与接入实战【免费下载链接】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信号发射管道emission pipeline是 PostHog Signals 产品把来自外部数据导入Zendesk 工单、GitHub issue、Linear 等和内部产品Conversations 工单的原始记录加工成可检索、可分组、可上报的 Signal 的核心后端组件。本文以 products/signals/backend/emission/AGENTS.md 为骨架结合源码实现完整讲解管道各阶段原理、直接来源与 steering 门控机制、注册表与抓取器设计并给出新增数据源与本地 fixture 冒烟测试的可运行实战方案。读完你可以独立理解现有 40 数据源的接入模式并照步骤为仓库添加一个新的 Signal 来源。一、发射出的 Signal 被如何使用description 字段的质量契约每一个被发射的 Signal 都会进入 Signals 工作流入口见 products/signals/backend/facade/api.py 中的emit_signal其description字段会被嵌入embedding用于语义搜索。来自不同来源、不同类型的 Signal 随后被合并进 Signal 分组signal groups再加工成 Signal 报告signal reports帮助用户发现产品中的问题。这意味着description字段必须为嵌入质量而写它应当以**与来源无关source-agnostic**的方式捕捉记录的含义让语义相似semantically similar的 Signal 无论来自哪个来源都能被良好地分到同一组。以 GitHub 为例github_issues.py 将title与body拼接为 descriptionZendesk 则以subjectdescription拼接见 zendesk_tickets.py。字段里混入的邮箱签名、法律免责声明、系统页脚等噪音需要在后面的总结阶段被剥离。二、共享管道Shared Pipeline的五阶段架构核心管道实现位于 pipeline.py 的run_signal_pipeline()它对来源完全中立source-agnostic任何来源都按同一套五阶段流水线处理Record fetcher记录抓取每个来源在自己的 config 上定义一个record_fetcher可调用对象负责把原始记录抓取为dict列表Emitter发射器把每条记录 dict 转换为SignalEmitterOutput若数据不足则返回None表示跳过Summarization总结可选阶段对超过阈值的长 description 通过 LLM 总结压缩Actionability filter可行动性过滤可选阶段通过 LLM 判断记录是否可行动actionable过滤掉噪音团队可通过SignalSourceConfig.config上的两个键对该闸门进行定向干预steeringEmission发射把存活下来的输出通过 products/signals/backend/facade/api.py 的emit_signal发射为 Signal。2.1 管道的具体运行逻辑源码级run_signal_pipeline的完整流程pipeline.py空记录短路无新记录时直接返回{status: success, reason: no_new_records, signals_emitted: 0}批量发射器执行build_emitter_outputs逐条调用 emitter单条抛异常只计数并跳过error_count 1日志中记录的是经过unloggable_fields脱敏的记录只有所有记录全部抛错时才抛出ApplicationError(non_retryableTrue)进入阶段打点每条输出都会打signal_data_source_entered事件capture_pipeline_stage事件名与属性见同一文件顶部总结阶段只有当 config 同时设置了summarization_prompt与description_summarization_threshold_chars时才触发超出阈值的 description 才送 LLM成功的打signal_data_source_summarized过滤阶段仅当 config 设置了actionability_prompt才触发被过滤的记录打signal_data_source_filtered携带steering_applied布尔值过滤后为空则返回no_actionable_records发射阶段_emit_signals用信号量限制并发EMIT_CONCURRENCY_LIMIT 50每条信号先估算 JSON 序列化后的字节数超过2MB 的 Temporal gRPC 载荷上限TEMPORAL_PAYLOAD_MAX_BYTES时先尝试丢弃extra重算仍超限则直接报错全部发射失败会抛RuntimeError让工作流重试。2.2 管道调优常量pipeline 中一组可直接阅读的调优常量pipeline.py常量默认值含义LLM_MODELclaude-sonnet-5可被环境变量SIGNAL_EMISSION_LLM_MODEL覆盖总结/行动性判断使用的模型LLM_CONCURRENCY_LIMIT20总结与行动性判断的并发 LLM 调用上限EMIT_CONCURRENCY_LIMIT50Signal 发射的并发上限LLM_MAX_ATTEMPTS3LLM 调用最大重试次数LLM_CALL_TIMEOUT_SECONDS120单次 LLM 调用超时LLM_RETRY_INITIAL_DELAY_SECONDS/LLM_RETRY_BACKOFF_COEFFICIENT5/2.0指数退避delay initial * coefficient^(attempt-1)RECORD_METADATA_MAX_CHARS2000注入行动性判断提示词的 record metadata 块长度上限LLM_MAX_OUTPUT_TOKENS8192刻意设高的输出 token 安全上限实际输出只有一句总结或一个单词总结失败时重试会遵循 Anthropic 要求 user/assistant 轮替的约束只有拿到过 assistant 文本才追加纠错消息所有重试耗尽后硬截断description 到阈值兜底output.description[:threshold]。行动性判断失败时则fail open——重试全部耗尽默认视为 actionablereturn True避免LLM 不可达就丢信号。三、直接来源Direct Sources与 direct_gate 门控错误追踪error tracking与健康检查health checks两类来源从不进入共享管道它们自行构建 description 后直接调用emit_signal。但它们仍然遵循团队的 steering 配置靠的是 direct_gate.py 中的闸门——emit_signal会对contracts.DIRECT_STEERABLE_SOURCES中列出的(source_product, source_type)组合应用该门控见 contracts.py 中的定义error_tracking/issue_created、error_tracking/issue_reopened、error_tracking/issue_spiking、health_checks/health_issue。这个门控有两条规则保证行为可预测没写 steering 就没有门控steering_filters_signal只读取文本形式的steering有意忽略已废弃的default_not_actionable因为直接来源的提示词本身不声明任何行动性标准启用它会几乎丢掉所有记录。团队没写任何 steering 时不产生 LLM 调用、不增加延迟、行为与今天完全一致团队规则是唯一标准DIRECT_SOURCE_ACTIONABILITY_PROMPT不像_prompts.py中的记录形态提示词那样自带判断标准它只声明团队偏好是唯一过滤器。所以写第一条规则只会过滤该规则描述的内容不会误伤其他。门控fail open异常时保留信号GATE_TIMEOUT_SECONDS 20秒超时被丢弃的信号会以signal_data_source_filtered带steering_applied打点与管道内过滤使用的阶段事件一致可区分被规则过滤与发射失败。两条纪律均由测试约束DIRECT_STEERABLE_SOURCES中的组合若同时被注册表registry提供每条记录会被判断两次因此两个集合必须不相交——tests/test_direct_gate.py 中的test_no_pipeline_source_is_listed_as_directly_steerable用isdisjoint断言守护这一点前端 agentRosterMeta.ts 中数据源的steerable标志必须镜像同一集合给没有门控支撑的来源展示 steering 表单等于存了没人读的文本。四、注册表Registrysource 与 config 的映射registry.py 以(source_type, schema_name)为键映射到各自的SignalSourceTableConfig。所有 emitter 在模块加载时自动注册_register_all_emitters()在文件末尾被调用。注册表键是普通字符串对外部来源使用ExternalDataSourceType的枚举值例如Zendesk内部来源使用自己的标识符例如conversationsInternalSourceType.CONVERSATIONS。一个细节GitHub 的 schema 行带仓库限定owner/repo.issues而 emitter 注册的是裸端点所以_registry_key()对 GitHub 会先经github_split_schema_name拆出 endpoint 再作键。get_signal_config()负责按规范化后的键查表。4.1 SignalSourceTableConfig数据源配置契约每个来源的 config 都是SignalSourceTableConfigPydantic 冻结模型registry.py字段含义如下字段说明source_product/source_type必须与SignalSourceConfig.SourceProduct/SourceType的选项一致emitter(team_id, record_dict) - SignalEmitterOutput \| None的纯函数record_fetcher来源自己定义抓取方式无默认值、必须显式指定partition_field用于时间窗过滤的字段如created_atfields要 SELECT 的列只含 emitter 与 extra 元数据所需where_clause可选过滤子句数据仓库源为 HogQLPostgres 源为 ORM 语法max_records每次同步最多处理的记录数默认1000partition_field_is_datetime_string为 True 时按字符串日期解析如 GitHub 的 JSON 字段first_sync_lookback_days首次同步回溯窗口默认7天actionability_prompt行动性判断提示词None表示所有记录都视为 actionableactionability_context_fields行动性闸门在 description 之外还需要读取的extra键按来源声明而非倾倒整个extraunloggable_fields日志前从记录中剥离的列承载超出 emitter 保留范围的更多身份信息summarization_prompt超阈值 description 的总结提示词None表示不做总结description_summarization_threshold_charsdescription 超长阈值须大于 0模型自带两条校验actionability_prompt/summarization_prompt必须包含{description}占位符summarization_prompt与description_summarization_threshold_chars必须成对出现要么都设要么都空。4.2 已注册的来源清单截至当前仓库从 registry.py 的_register_all_emitters()可以看到按 Tier/记录形态组织的完整清单Tier-1 工单/客服record kind: ticketZendesk、Freshdesk、Freshservice、Front、Gorgias、Kustomer、Dixa、Plain、Intercom、HubSpotTier-1 issue 追踪器record kind: issueGitHub、Linear、Jira、pganalyze、GitLab、Gitea、ShortcutTier-1 错误追踪record kind: issueSentry、Rollbar、Bugsnag、Honeybadger、RaygunTier-2 安全扫描器record kind: scanner_findingSnyk、SonarQube、Semgrep、Rapid7 InsightVMTier-3 产品反馈/功能请求record kind: feedbackFeaturebase、Frill、Aha、Uservoice、Productboard、Canny、AskNicely、RetentlyTier-3 应用商店评价record kind: reviewAppfigures、Appfollow、Judge.me搜索分析record kind: search_opportunityGoogle Search Console内部产品Conversationsconversations/tickets。五、记录抓取器Record Fetchers的两种实现每个来源通过 config 上的record_fetcher定义抓取方式当前仓库有两种5.1 数据仓库抓取器HogQLfetchers/data_warehouse.py 的data_warehouse_record_fetcher通过 HogQL 查询仓库表。运行时上下文 dict 传入table_name与last_synced_at持续同步last_synced_at存在时生成partition_expr {last_synced_at}条件占位符用ast.Constant绑定避免注入首次同步无last_synced_at时用now() - interval {first_sync_lookback_days} day限定回看窗口若partition_field_is_datetime_string为 True分区表达式会包一层parseDateTimeBestEffort(...)表名按点分段逐个escape_hogql_identifier转义escape_table_name因为部分来源把客户文本放进表键如 GitHub 仓库名带连字符查询走execute_hogql_query(query_typeEmitSignalsNewRecords, bypass_warehouse_access_controlTrue)——内部抓取器无用户身份需绕过仓库访问控制才能读来源仓库表查询失败会重新抛出而不是吞掉避免在发射信号的时间轴上制造永久空洞。5.2 Conversations 抓取器Django ORMfetchers/conversations.py 的conversations_ticket_fetcher用 Django ORM 查询 Postgres 中的工单Ticket与评论Comment两条时间规则TICKET_QUIET_PERIOD_HOURS 1线程必须安静满 1 小时以最后一条消息为准回退到created_at才做快照——因为支持线程里的决定性细节复现步骤、报错文本、升级通常出现在回复而非开场消息中TICKET_RESNAPSHOT_MIN_INTERVAL_HOURS 24同一工单两次快照的最短间隔活跃线程最多每天重新快照一次避免每个安静间隙都发一条信号评论查询排除私密便签与 AI 草稿is_private把图片 URL 从rich_content中抽出作为image_attachments乐观记录发射抓取返回前就批量 upsertSignalEmissionRecordupdate_conflictsTrue唯一键team source_product source_type source_idemitted_at同时充当下次重快照的门槛。注意 emit_signals.py 与 conversations_coordinator.py 都特意在抓取之前读取 source_config因为一旦 fetcher 返回就乐观记录了发射状态之后若 Temporal 重试会永久跳过这批。六、两类触发链路数据导入工作流与 Conversations 定时调度6.1 数据导入来源Zendesk、GitHub、Linear 等由数据导入工作流触发两级工作流父工作流external_data_job.py 完成数据导入后若该来源启用了信号发射则派生 emit-signals 子工作流子工作流emit_signals.py 中名为emit-data-import-signals的EmitDataImportSignalsWorkflow执行 activityemit_data_import_signals_activity先按(source_type, schema_name)查注册表未注册则跳过no_config_registered随后加载ExternalDataSchema与Team构造 fetcher 上下文table_name用get_data_warehouse_table_name规范化、last_synced_at、日志属性调用config.record_fetcher后进入共享管道。Activity 配置emit_signals.pystart_to_close_timeout60 分钟、heartbeat_timeout5 分钟、RetryPolicy(maximum_attempts3)。6.2 Conversations 来源Temporal 每小时调度由 Temporal schedule 每小时触发一次两级工作流协调器工作流conversations_coordinator.pyconversations-signals-coordinator先查询启用了 conversations 信号且通过 AI 数据审批的团队列表get_conversations_signals_enabled_teams_activity然后分批派生每团队子工作流每批最多DEFAULT_MAX_CONCURRENT_TEAMS 50个并发为避免超长 rollout 撑爆 Temporal 历史用CoordinatorStateremaining team ids 成功/失败/发射计数配合continue_as_new续跑每团队工作流emit-conversations-signals执行 activity抓取合格工单1 小时安静、未解决、尚未发射带完整消息线程再跑共享管道。七、门控Gating所有来源的统一准入条件所有来源无论走共享管道还是直接发射都被两道门拦住AI 数据审批组织必须满足organization.is_ai_data_processing_approvedConversations 协调器查询团队列表时也叠加了这一条件信号启用理应要求 AI 审批来源启用开关存在匹配source_product/source_type且enabledTrue的SignalSourceConfig行。用户通过 Inbox 的 Sources 弹窗Inbox Sources modal启用来源。八、Steering团队对行动性闸门的定向干预SignalSourceConfig.config上有两个公开键定义见 contracts.py键类型含义steeringstring上限STEERING_MAX_LENGTH 2000字符团队用自然语言描述的偏好什么重要、什么跳过、什么超范围default_not_actionablebool翻转闸门姿态从默认保留keep everything except what rules exclude变为只保留明确符合的only keep what clearly qualifies关键设计steering.py团队提供规则而非提示词注入发生在提示词模板层.format(description...)之前且把{/}转义为{{/}}恶意输入格式串语法、游离花括号无法让后续format抛异常从而触发 fail-open 路径姿态锚点所有规范行动性提示词都保留When in doubt, classify as ACTIONABLE这一行_POSTURE_MARKER注入与姿态翻转都锚定该行default_not_actionable会把该行替换为 allowlist 姿态行只保留明确匹配 ACTIONABLE 标准的记录元数据块行动性判断发生在信号存在之前闸门只能看到 description——除非记录在别处携带了判决依据。来源在actionability_context_fields声明这些extra键它们会以record_metadata块的形式附进 description见 pipeline.py 的_declared_context与check_actionability块长度受RECORD_METADATA_MAX_CHARS限制。被 steering 的团队看到整个extra声明为空的来源提示词保持逐字节不变容错解析steering_from_config对畸形值防御式降级非 Mapping 返回空 steering、非字符串截断、default_not_actionable严格is True保证 API/MCP 写入的脏 JSON 只退化为规范行为而不破坏发射。以 github_issues.py 为例它声明了actionability_context_fields(author_login, author_association)——因为谁提交的报告决定了维护者 bug 报告与路人报告的权重差异triage 不能没有它同时该来源的unloggable_fields(user,)发射器只从嵌套user对象中提起login句柄并丢弃头像/API URL 等一堆 URL。生命周期的遥测事件facade/api.py 的_TELEMETRY_EXCLUDED_EXTRA_KEYS同样排除身份键否则extra上每个顶层标量都会被复制进发射事件。九、新增一个数据源的分步指南以文档给出的 Jira 为例对应源码中已有 jira_issues.py 实现可作参照第 1 步创建 emitter 模块在本目录products/signals/backend/emission/下新建文件如jira_issues.py参照 zendesk_tickets.py、github_issues.py、conversations_tickets.py 的模式定义要查询的字段REQUIRED_FIELDS 透传/附加元数据字段写一个纯函数 emitter把记录 dict 转为SignalEmitterOutputsource_product、source_type、source_id、description、weight、extra数据不足返回None定义record_fetcher仓库来源用data_warehouse_record_fetcher其他来源自写新 fetcher可选地定义 LLM 行动性提示词与/或带阈值的总结提示词可复用 _prompts.py 中按记录形态分组的共享提示词——工单、issue、错误、扫描发现、反馈、评价六种形态把最终 config 导出为模块级常量如JIRA_ISSUES_CONFIG。PII 红线除非严格需要用于在来源系统中定位实体否则避免查询 PII 字段用户 ID、邮箱、姓名、组织 ID 等优先使用不透明的记录 ID 与 URL。唯一被刻意允许的例外是记录作者github_issues.py携带author_login与author_association因为公开仓库上谁提交的报告是区分维护者 bug 报告与路人报告的关键triage 缺了它无法权衡。其他 issue 追踪器来源应沿用同一对字段且只保留句柄与关系——不要邮箱、真名或嵌套 user 对象的其余部分。如果一个来源 SELECT 了比它保留的更多身份的列必须把这些列声明进unloggable_fields共享管道在 emitter 抛异常记录日志时会用redacted_record剥掉它们registry.py身份键同时被facade/api.py的_TELEMETRY_EXCLUDED_EXTRA_KEYS排除在发射遥测之外。第 2 步在 registry.py 中注册在 registry.py 的_register_all_emitters()内导入 config 并调用register_signal_source(...)外部来源用ExternalDataSourceType的值作为 source type内部来源用描述性字符串标识符参照InternalSourceType.CONVERSATIONS conversations。第 3 步在 tests/ 中写测试emitter 测试test_source.py覆盖合法记录、缺失/空必需字段参数化、extra 字段提取在 tests/conftest.py 中追加贴近现实的 mock 记录与 pytest fixture。运行测试pytest products/signals/backend/emission/tests/仓库的 tests/ 目录里已有test_zendesk_tickets.py、test_github_issues.py、test_conversations_fetcher.py、test_steering.py、test_direct_gate.py、test_registry.py、test_schema_validation.py等 16 个测试文件可作为模式参考。十、本地 fixture 冒烟测试不走真实导入跑通全管道要在不跑真实数据导入、不填充仓库表的情况下演练完整管道emitter → summarization → actionability →emit_signal使用emit_signals_from_fixture管理命令它从 products/signals/eval/fixtures/ 加载脱敏 fixture 记录直接喂给run_signal_pipeline完全绕过data_warehouse_record_fetcher。# 用 1-2 条记录做冒烟测试成本低每条记录约 1-2 次 LLM 调用 DEBUG1 ./manage.py emit_signals_from_fixture --type zendesk --team-id 1 --limit 1 DEBUG1 ./manage.py emit_signals_from_fixture --type github --team-id 1 --limit 2 DEBUG1 ./manage.py emit_signals_from_fixture --type linear --team-id 1 DEBUG1 ./manage.py emit_signals_from_fixture --type conversations --team-id 1 --limit 2 # 覆盖 fixture 路径 DEBUG1 ./manage.py emit_signals_from_fixture --type zendesk --team-id 1 --fixture path/to/custom.json参数说明--type接受zendesk、github、linear或conversations映射到 registry.py 中对应的自动注册 config命令要求DEBUGTrue仅用于本地迭代如 pipeline.py 的_safe_heartbeat所示管道既能在 Temporal activity 内运行也能独立于 activity 上下文以管理命令方式运行steering 同样生效fixture 运行以及emit_signals_from_llm会像生产环境一样读取团队的SignalSourceConfig.configsteering 键所以设置了steering或default_not_actionable的团队可能过滤掉普通运行会保留的记录要拿到无 steering 的基线先清掉该来源 config 行上的这两个键。十一、维护约定与阅读延伸文档最后明确了维护契约当管道架构、注册表模式或接入约定发生重大变化时需要同步更新 AGENTS.md 以反映新现实。与之配套的可读材料还包括系统架构总览products/signals/ARCHITECTURE.md其中也描述了steering/default_not_actionable两键对管道来源与DIRECT_STEERABLE_SOURCES直接来源的适用范围信号载荷契约SignalSourceConfig两键与DIRECT_STEERABLE_SOURCES的定义见 products/signals/backend/contracts.py所有发射载荷在发射边界按这些模型校验extraforbid未知字段直接拒绝按记录形态共享的提示词模板products/signals/backend/emission/_prompts.py。整个发射管道的设计取向可以概括为三点description 为嵌入而写来源无关的语义纯度、闸门可控且 fail opensteering 注入有防破坏保证、LLM 不可达不丢真信号、身份最小化PII 列不查询、不多留、不落入日志与遥测。【免费下载链接】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),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门