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拼接为 description;Zendesk 则以subject+description拼接(见 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上的两个键对该闸门进行定向干预(steering); - Emission(发射):把存活下来的输出通过 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_retryable=True); - 进入阶段打点:每条输出都会打
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_MODEL | claude-sonnet-5(可被环境变量SIGNAL_EMISSION_LLM_MODEL覆盖) | 总结/行动性判断使用的模型 |
LLM_CONCURRENCY_LIMIT | 20 | 总结与行动性判断的并发 LLM 调用上限 |
EMIT_CONCURRENCY_LIMIT | 50 | Signal 发射的并发上限 |
LLM_MAX_ATTEMPTS | 3 | LLM 调用最大重试次数 |
LLM_CALL_TIMEOUT_SECONDS | 120 | 单次 LLM 调用超时 |
LLM_RETRY_INITIAL_DELAY_SECONDS/LLM_RETRY_BACKOFF_COEFFICIENT | 5/2.0 | 指数退避delay = initial * coefficient^(attempt-1) |
RECORD_METADATA_MAX_CHARS | 2000 | 注入行动性判断提示词的 record metadata 块长度上限 |
LLM_MAX_OUTPUT_TOKENS | 8192 | 刻意设高的输出 token 安全上限(实际输出只有一句总结或一个单词) |
总结失败时,重试会遵循 Anthropic 要求 user/assistant 轮替的约束,只有拿到过 assistant 文本才追加纠错消息;所有重试耗尽后硬截断description 到阈值兜底(output.description[:threshold])。行动性判断失败时则fail open——重试全部耗尽默认视为 actionable(return 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 表单,等于存了没人读的文本。
四、注册表(Registry):source 与 config 的映射
registry.py 以(source_type, schema_name)为键映射到各自的SignalSourceTableConfig。所有 emitter 在模块加载时自动注册(_register_all_emitters()在文件末尾被调用)。注册表键是普通字符串对:
- 外部来源使用
ExternalDataSourceType的枚举值,例如"Zendesk"; - 内部来源使用自己的标识符,例如
"conversations"(InternalSourceType.CONVERSATIONS)。
一个细节:GitHub 的 schema 行带仓库限定(owner/repo.issues),而 emitter 注册的是裸端点,所以_registry_key()对 GitHub 会先经github_split_schema_name拆出 endpoint 再作键。get_signal_config()负责按规范化后的键查表。
4.1 SignalSourceTableConfig:数据源配置契约
每个来源的 config 都是SignalSourceTableConfig(Pydantic 冻结模型,registry.py),字段含义如下:
| 字段 | 说明 |
|---|---|
source_product/source_type | 必须与SignalSourceConfig.SourceProduct/SourceType的选项一致 |
emitter | (team_id, record_dict) -> SignalEmitterOutput \| None的纯函数 |
record_fetcher | 来源自己定义抓取方式,无默认值、必须显式指定 |
partition_field | 用于时间窗过滤的字段(如created_at) |
fields | 要 SELECT 的列,只含 emitter 与 extra 元数据所需 |
where_clause | 可选过滤子句(数据仓库源为 HogQL,Postgres 源为 ORM 语法) |
max_records | 每次同步最多处理的记录数,默认1000 |
partition_field_is_datetime_string | 为 True 时按字符串日期解析(如 GitHub 的 JSON 字段) |
first_sync_lookback_days | 首次同步回溯窗口,默认7天 |
actionability_prompt | 行动性判断提示词;None表示所有记录都视为 actionable |
actionability_context_fields | 行动性闸门在 description 之外还需要读取的extra键(按来源声明,而非倾倒整个extra) |
unloggable_fields | 日志前从记录中剥离的列(承载超出 emitter 保留范围的更多身份信息) |
summarization_prompt | 超阈值 description 的总结提示词;None表示不做总结 |
description_summarization_threshold_chars | description 超长阈值(须大于 0) |
模型自带两条校验:actionability_prompt/summarization_prompt必须包含{description}占位符;summarization_prompt与description_summarization_threshold_chars必须成对出现(要么都设,要么都空)。
4.2 已注册的来源清单(截至当前仓库)
从 registry.py 的_register_all_emitters()可以看到按 Tier/记录形态组织的完整清单:
- Tier-1 工单/客服(record kind: ticket):Zendesk、Freshdesk、Freshservice、Front、Gorgias、Kustomer、Dixa、Plain、Intercom、HubSpot;
- Tier-1 issue 追踪器(record kind: issue):GitHub、Linear、Jira、pganalyze、GitLab、Gitea、Shortcut;
- Tier-1 错误追踪(record kind: issue):Sentry、Rollbar、Bugsnag、Honeybadger、Raygun;
- Tier-2 安全扫描器(record kind: scanner_finding):Snyk、SonarQube、Semgrep、Rapid7 InsightVM;
- Tier-3 产品反馈/功能请求(record kind: feedback):Featurebase、Frill、Aha、Uservoice、Productboard、Canny、AskNicely、Retently;
- Tier-3 应用商店评价(record kind: review):Appfigures、Appfollow、Judge.me;
- 搜索分析(record kind: search_opportunity):Google Search Console;
- 内部产品:Conversations(
conversations/tickets)。
五、记录抓取器(Record Fetchers)的两种实现
每个来源通过 config 上的record_fetcher定义抓取方式,当前仓库有两种:
5.1 数据仓库抓取器(HogQL)
fetchers/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_type="EmitSignalsNewRecords", bypass_warehouse_access_control=True)——内部抓取器无用户身份,需绕过仓库访问控制才能读来源仓库表; - 查询失败会重新抛出而不是吞掉,避免在发射信号的时间轴上制造永久空洞。
5.2 Conversations 抓取器(Django ORM)
fetchers/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; - 乐观记录发射:抓取返回前就批量 upsert
SignalEmissionRecord(update_conflicts=True,唯一键team + source_product + source_type + source_id),emitted_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执行 activity:emit_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.py):start_to_close_timeout=60 分钟、heartbeat_timeout=5 分钟、RetryPolicy(maximum_attempts=3)。
6.2 Conversations 来源(Temporal 每小时调度)
由 Temporal schedule 每小时触发一次,两级工作流:
- 协调器工作流conversations_coordinator.py(
conversations-signals-coordinator):先查询启用了 conversations 信号且通过 AI 数据审批的团队列表(get_conversations_signals_enabled_teams_activity),然后分批派生每团队子工作流,每批最多DEFAULT_MAX_CONCURRENT_TEAMS = 50个并发;为避免超长 rollout 撑爆 Temporal 历史,用CoordinatorState(remaining team ids + 成功/失败/发射计数)配合continue_as_new续跑; - 每团队工作流(
emit-conversations-signals)执行 activity:抓取合格工单(>1 小时安静、未解决、尚未发射,带完整消息线程),再跑共享管道。
七、门控(Gating):所有来源的统一准入条件
所有来源(无论走共享管道还是直接发射)都被两道门拦住:
- AI 数据审批:组织必须满足
organization.is_ai_data_processing_approved(Conversations 协调器查询团队列表时也叠加了这一条件,"信号启用理应要求 AI 审批"); - 来源启用开关:存在匹配
source_product/source_type且enabled=True的SignalSourceConfig行。
用户通过 Inbox 的 Sources 弹窗(Inbox Sources modal)启用来源。
八、Steering:团队对行动性闸门的定向干预
SignalSourceConfig.config上有两个公开键(定义见 contracts.py):
| 键 | 类型 | 含义 |
|---|---|---|
steering | string(上限STEERING_MAX_LENGTH = 2000字符) | 团队用自然语言描述的偏好:什么重要、什么跳过、什么超范围 |
default_not_actionable | bool | 翻转闸门姿态:从"默认保留(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 转为
SignalEmitterOutput(source_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 调用) DEBUG=1 ./manage.py emit_signals_from_fixture --type zendesk --team-id 1 --limit 1 DEBUG=1 ./manage.py emit_signals_from_fixture --type github --team-id 1 --limit 2 DEBUG=1 ./manage.py emit_signals_from_fixture --type linear --team-id 1 DEBUG=1 ./manage.py emit_signals_from_fixture --type conversations --team-id 1 --limit 2 # 覆盖 fixture 路径 DEBUG=1 ./manage.py emit_signals_from_fixture --type zendesk --team-id 1 --fixture path/to/custom.json参数说明:
--type接受zendesk、github、linear或conversations,映射到 registry.py 中对应的自动注册 config;- 命令要求
DEBUG=True,仅用于本地迭代(如 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,所有发射载荷在发射边界按这些模型校验(extra="forbid",未知字段直接拒绝); - 按记录形态共享的提示词模板:products/signals/backend/emission/_prompts.py。
整个发射管道的设计取向可以概括为三点:description 为嵌入而写(来源无关的语义纯度)、闸门可控且 fail open(steering 注入有防破坏保证、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),仅供参考