Sentry Workflow Engine 执行链路深度解析:从 DataPacket 到 Action 的完整数据流
【免费下载链接】sentryDeveloper-first error tracking and performance monitoring项目地址: https://gitcode.com/GitHub_Trending/sen/sentry
导读
本文以 src/sentry/workflow_engine/docs/execution.md 为骨架,结合 Sentry 源码,系统讲解 Workflow Engine 的两条执行流水线:由 Detector 驱动的检测流水线(Detection Pipeline)与由 Issue 事件/活动驱动的工作流流水线(Workflow Pipeline),以及贯穿其中的 Kafka 异步边界、Redis 延迟缓冲、Snuba 批量查询与 Action 分发机制。读完本文,你将掌握DataPacket从产生到发布 Issue Platform 的完整路径、StatefulDetectorHandler的有状态阈值判定语义、process_workflows的 WHEN/IF 求值流程、延迟工作流的缓冲—调度—评估三步曲,以及 Action 重复频率控制与去重的底层实现。
一、两条流水线:检测与工作流通过 Issue Platform 解耦
Workflow Engine 的核心设计是检测(Detection)与工作流处理(Workflow Processing)之间不存在直接函数调用,二者通过 Issue Platform 连接:
- Detector 的输出被发布到 Issue Platform(Kafka 主题);
- 由 Issue Platform 消化生成的 issue 事件,随后通过post-processing或issue activity进入工作流处理。
这种解耦带来的直接结果是:Kafka 边界意味着"检测评估成功"并不等价于"group 创建完成、工作流已处理、Action 已派发"。检测与动作执行之间存在天然异步,这也是本文后面"失败与一致性模型"一节讨论的主题。
另外,普通错误事件可以不经过任何 Workflow Engine 检测器,直接在 post-processing 阶段进入工作流流水线——这正是错误工作流(error workflows)与 issue-stream 工作流能共用同一套工作流管线的原因。相关的持久化关系请参阅>DataPacket(source_id=str(source.id), packet=product_payload)
要点:
packet载荷保持产品特定(product-specific),Workflow Engine 不解析其内部结构;- 生产方同时提供source type(
query_type),用于解析source_id到底对应哪个DataSource。
2. 解析 Detector
process_data_packet是通用的入口函数,它把两个预处理步骤串联起来,并定义了可用于监控引擎健康度的指标漏斗(workflow_engine.process_data_sources→workflow_engine.process_detectors→workflow_engine.process_workflows):
process_data_packet └─ process_data_source # 解析 (source type, source_id) -> Detector 列表 └─ process_detectors # 逐个评估 detector,产出 DetectorEvaluationprocess_data_source的解析链路为:
(source type, source_id) -> DataSource -> DataSourceDetector -> Detector rows with enabled=True through the default manager源码细节:
bulk_fetch_enabled_detectors先通过get_detectors_by_data_source走缓存,再在缓存层之外过滤d.enabled(return [d for d in detectors if d.enabled]),避免缓存层承担过滤职责;- 查找使用 detector-by-source 缓存(见"缓存"一节),并预加载触发条件(
prefetch_related("workflow_condition_group__conditions"),见 caches/detector.py); - 缓存失效连接到 source、detector 与 mapping 的更新(对应 receivers 目录下的各 receiver);
- 注意:默认 manager 会排除 pending-deletion 与 deletion-in-progress 的行,但该路径并不会排除所有非 active 状态(例如计划级别的
DISABLED),这是有意的取舍。
并非所有生产方都走这个通用查找。Uptime、processing errors、preprod 等路径可以通过产品特定数据直接识别 detector,然后调用process_detectors。detector 选定之后,handler 的契约是相同的。
3. 逐个评估 Detector
process_detectors的核心逻辑:
for detector in detectors: handler = detector.detector_handler if not handler: continue detector_results = handler.evaluate(data_packet) for result in detector_results.values(): if result.result is not None: create_issue_platform_payload(result, detector.type)其中detector.detector_handler来自 detector 注册的GroupType.detector_settings。一个 packet 可能产生:
- 没有评估结果(没有条件被触发);
- 一个结果(未分组 detector,
DetectorGroupKey为None); - 多个结果(按
DetectorGroupKey分组)。
Handler 返回一个dict[DetectorGroupKey, DetectorEvaluation],每个 evaluation 可以包含:
- 一个
IssueOccurrence(触发); - 一个
StatusChangeMessage(恢复/解决); - 或
None(无输出)。
process_detectors只发布非 null 结果,并会按结果类型打workflow_engine.process_detector.triggered/.resolved指标。
4. 有状态 Detector 编排(Stateful Detector Orchestration)
大多数基于阈值的 detector 继承StatefulDetectorHandler,其流程如下:
批量状态 I/O:handler 对分组值(grouped values)执行批量状态读写。具体的事务与 Redis pipeline 行为由同模块的DetectorStateManager实现,关键点如下:
get_state_data先批量拉取DetectorState(PostgreSQL)行,再通过单个 Redis pipeline一次性GET所有 dedupe 与 counter key(见 stateful.py 中的get_redis_keys_for_group_keys/bulk_get_redis_values);- 提交时同样用 pipeline 批量
SET(带ex=REDIS_TTL)或DELETE,对DetectorState则用bulk_create/bulk_update(_bulk_commit_redis_state/_bulk_commit_detector_state)。
关键语义(务必理解):
- Dedupe 值必须是正整数,且按事件顺序递增。只要 Redis watermark 存在,等于或更小的值都会被忽略;watermark 缺失时默认 0,7 天后过期(源码中
REDIS_TTL = int(timedelta(days=7).total_seconds())); - 状态按detector group key 相互独立;
- 达到更高优先级的同时,也会递增所有适用的低优先级计数器(见
_increment_detector_thresholds:对level <= new_priority且非 OK 的 level 全部 +1); - 阈值描述的是:按优先级计数器规则,需要多少次合格评估才会发生状态迁移;
- 迁移到非 OK 优先级 → 创建 occurrence;
- 迁移到
OK→ 创建 status-change message; - 无状态迁移 → 不产生任何 Issue Platform 消息。
恢复(Recovery)语义:并非"所有非 OK 条件都失败"就意味着恢复——因为一个未被触发的条件组不会产生任何迁移。恢复要求存在被触发(triggered)的条件组,且其选中的优先级保持OK。通常这来自一个结果为OK的通过条件;但一个空或通过的NONE条件组也可以被触发(不携带优先级结果),此时使用 stateful handler 的默认OK优先级。缺失触发条件组是无效配置,不会产生任何迁移。
自定义 detector 可以继承更小的BaseDetectorHandler,或直接实现DetectorHandler接口,但那样就需要自己承担上述大部分编排逻辑。
5. 发布到 Issue Platform
process_detectors将结果交给create_issue_platform_payload与produce_occurrence_to_kafka:
- 载荷类型(
PayloadType.OCCURRENCE/PayloadType.STATUS_CHANGE)区分 occurrence 与状态变更; - Issue Platform 消化后创建或更新 group;occurrence evidence 中的detector ID允许消化逻辑创建
DetectorGroup关联; - 恢复(resolution)使用稳定的 issue fingerprint找到同一个 group。
相关实现还有associate_new_group_with_detector/ensure_association_with_detector(见 processors/detector.py):当 detector 已被删除时会创建detector_id=None的DetectorGroup占位,以显式表达"曾关联到一个已不存在的 detector"。
三、工作流入向(Workflow Ingress)
事件后处理(Event post-processing)
process_workflow_engine只对非 reprocessed 的 job、且事件的 group 当前处于未解决(unresolved)状态时,调度process_workflows_event任务。任务通过 tasks/utils.py 从 EventStore 与 Issue Platform 数据重建WorkflowEventData。基于 resolution 的处理则改由受支持的 issue activity 进入。
Issue activities
Issue activity 通过invoke_workflow_activity_handlers调用处理器。已注册的 handler 会收到可选的 detector ID;任何要调度process_workflow_activity的 handler,都必须先解析出具体的 detector,并把必需的 detector ID 连同 activity 和 group一起传给任务(任务签名process_workflow_activity(activity_id, group_id, detector_id),见 tasks/workflows.py)。
WorkflowEventData 的内容与两条路径的差异
事件路径与活动路径都使用WorkflowEventData,它包含:
- 一个
GroupEvent或Activity; - 对应的 issue
Group; - 可选的 group 状态与升级(escalation)信息;
- 可选的工作流环境(workflow environment);
- 一个每事件本地缓存(per-event local cache),用于重复查询去重。
但两条路径调度并不完全相同:
- 活动处理没有工作流环境,因此只考虑全局(global,
environment=None)工作流; - 活动不进入批量的条件数据获取:一个 activity 的 WHEN 或 IF 结果如果还需要慢查询(slow query),不会继续进入延迟处理(源码中
evaluate_workflow_triggers与evaluate_workflows_action_filters均对Activity分支直接metrics_incr("process_workflows.enqueue_workflow.activity")并跳过入队)。
四、工作流流水线(Workflow Pipeline)
process_workflows是工作流任务内部的主同步编排器:
process_workflows还会以WorkflowEvaluationOutcome(NO_DETECTOR/ENVIRONMENT_NOT_FOUND/NO_WORKFLOWS/COMPLETED)返回结果,方便上层统计与观测。
1. 解析事件 Detector
get_detectors_for_event_data构建EventDetectors对象:
- 对于带 occurrence 的
GroupEvent,event_detector只从 occurrence 的detector_idevidence 解析(issue_occurrence.evidence_data.get("detector_id"));ID 缺失或被删除 → 没有事件 detector; - 对于没有 occurrence 的普通错误事件,使用项目默认错误 detector(
Detector.get_error_detector_for_project); issue_stream_detectors包含项目的 issue-stream detector,以及(启用时)组织的all-projects detector(按in_rollout_group("workflow_engine.all_projects_detectors.rollout-rate", ...)灰度开启);preferred_detector优先取事件 detector,否则取第一个 issue-stream detector。
EventDetectors是一个 frozen dataclass:至少要有 1 个 detector,否则__post_init__抛ValueError,get_detectors_for_event_data捕获后返回None(见 processors/detector.py)。preferred detector 之后会提供给需要 detector 上下文的 action。由于 detector 删除与较老的 group 可能遗留缺失的关联,所有查找路径都被设计为可安全降级(degrade safely)。
2. 解析环境与工作流
处理器先解析事件环境,再由get_workflows_by_detectors加载与任一选中 detector 相连的工作流。合格的工作流必须:
enabled=True;- 是全局的(
environment=None)或连接到事件的 environment。
缓存按detector + environment键控,关系 receiver 在提交(committed)后使其失效。源码中还保留了一条 feature-flag 关闭缓存时的直接 DB 查询回退路径_get_associated_workflows(Workflow.objects.filter(environment_filter, detectorworkflow__detector_id__in=..., enabled=True))。
3. 评估 WHEN 条件
evaluate_workflow_triggers对每个工作流求值其 WHEN 条件组:
- 没有 WHEN 组的工作流直接通过;
- WHEN 结论为 false → 停止处理并记录 metrics/context,返回的 per-workflow 结果 map可以保持为空;
- WHEN 有**慢条件(slow conditions)**且事件是
GroupEvent→ 构造DelayedWorkflowItem等待缓冲;若是Activity→ 不入队(见上文差异); - 为减少查询,所有 WHEN 组的 DataCondition 与 DataConditionGroup 会批量获取(
get_data_conditions_for_group.batch与get_many_from_cache)。
4. 评估 IF 条件组
evaluate_workflows_action_filters对 WHEN 通过或仍处于慢求值待定状态的工作流,从action-filter 缓存加载 IF 组:
- 快速的 IF 结果与剩余的慢 IF 组会与待定的 WHEN 结果一起保留,以便合并缓冲;
- 通过的条件组贡献其挂接的 action——一个工作流可以从多个条件组中挑选 action;
- 没有 IF 组的工作流不挑选任何 action;
- 慢条件入队的合并规则:同一条
DelayedWorkflowItem上追加delayed_if_group_ids;已通过的 IF 组记录到passing_if_group_ids(它们依赖 WHEN 延迟求值通过后才触发)。
共享求值、fast/slow 拆分、空条件组行为、taint 契约以及 Detector 差异,在 conditions.md 中有统一说明。
五、延迟工作流处理(Delayed Workflow Processing)
当 WHEN 或 IF 条件需要 Snuba 慢查询(如事件频率、百分比比较条件)时,工作流进入延迟路径,跨工作流、跨 issue group 批量执行昂贵的条件:
缓冲(Buffering)
DelayedWorkflowClient将 project ID 存入分片的有序集合(sorted set),每个 project 的工作存入一个 hash:
- hash 字段以
workflow_id:group_id:when_dcg_id:if_dcg_ids:passing_dcg_ids为 key(见DelayedWorkflowItem.buffer_key); - 字段值保留事件或 occurrence ID 与原始评估时间戳(
buffer_value序列化event_id/occurrence_id/timestamp); - 对同一字段的重复入队会覆盖旧值——合并等价待处理工作,而不是保留每个事件;
- 底层为分片 buffer:
_BUFFER_KEY = "workflow_engine_delayed_processing_buffer"、_BUFFER_SHARDS = 8,project 被随机放入某个分片;rollout 由 optiondelayed_workflow.rollout控制。
调度(Scheduling)
schedule_delayed_workflows任务在全局调度锁内执行process_buffered_workflows:
- 通过
ProjectChooser把 project 划分到 cohort(hashlib.sha256(project_id).digest()[0] % num_cohorts),按"must-process(超过 1 分钟)→ 最陈旧 cohort"的策略选出待处理 project(target_max_age = 1 minute,min_scheduling_age默认 50 秒,见 optionworkflow_engine.schedule.min_cohort_scheduling_age_seconds); - 大 project hash 会按
delayed_processing.batch_size分批:每批以uuid4()作为 batch key 复制到新 hash、从原 hash 删除,再调度process_delayed_workflows(project_id, batch_key)——异步任务消息里只带标识符,不携带大载荷(源码注释明确说明:10k 条项目无法作为任务参数传递,Redis 也不保持 hash key 顺序,无法分页); - 处理完成后用
mark_project_ids_as_processed按"小于等于该时间戳"条件删除已处理成员。
评估(Evaluation)
process_delayed_workflows的执行步骤:
- 解析缓冲项(
EventRedisData,容忍坏条目时continue_on_error=True); - 丢弃对已删除 workflow 的引用(
filter_by_workflow_ids); - 加载剩余条件;
- 把等价条件折叠为批量查询(
UniqueConditionQuery按 handler + interval + environment_id + comparison_interval + frozen filters 键控;百分比比较条件会产生两条唯一查询); - 按 issue group 求值查询结果;
- 将延迟结果与缓冲 key 中保留的快速结果合并;
- 从 nodestore 批量重建代表性 group 事件(
bulk_fetch_events,每次 100 条EVENT_LIMIT,带should_retry_fetch重试策略); - 通过正常 action 路径派发通过的 action 组(
fire_actions_for_groups,fire history 标记is_delayed=True); - 删除成功处理的 Redis 项(
cleanup_redis_buffer)。
容错细节:
- Snuba 限流(
RateLimitExceeded)由任务重试;最后一次尝试仍遇限流时记录日志并按 false 处理,避免整个任务失败; - 某个 group 缺少查询结果 → 产生tainted 的非触发评估(不 firing);
- 慢条件没有注册查询 handler时不生成查询,处理可以正常完成并移除缓冲项(不构造 tainted 结果);
- 格式损坏的 hash 条目会被排除在解析出的事件 key 集合之外,不会被正常成功清理删除,会一直保留到过期或人工清理。
六、Action 选择与分发(Action Selection and Dispatch)
processors/action.py施加两道不同的控制:
重复频率(Repeat frequency)
WorkflowActionGroupStatus记录每个(workflow, action, issue group)元组最近一次通过频率检查的时间。它在事件级去重与任务执行之前被更新:
- 频率间隔来自
Workflow.config.get("frequency", 0)(单位分钟,见process_workflow_action_group_statuses中的frequency * timedelta(minutes=1)); - 当
now - date_updated未超过频率间隔时,该元组被抑制(suppressed)——即使之后有等价的 action 被去重,或外部 handler 执行失败,该元组仍被抑制; - 缺失状态(missing status)会以
ON CONFLICT (workflow_id, action_id, group_id) DO NOTHING批量插入;插入遇IntegrityError时切换为带 EXISTS 检查的 SQL 重试,以跳过引用已删除的孤儿行。
事件级去重(Event-level deduplication)
get_unique_active_actions两步走:
- 先剔除
status != ACTIVE的 action; - 再对每个
Action.get_dedup_key()只保留一个具体 action。
这是一个实现级 key,并不保证所有语义等价的外部副作用都能被识别。调用的归属(attribution)来自该 key 保留的那个具体 action,不暗示工作流执行顺序。
历史与任务分发(History and task dispatch)
处理器在fire_actions之前更新频率状态并创建WorkflowFireHistory行,然后调度trigger_action:
- fire-history 行只描述"已调度",不代表外部投递成功——它们可以存在于之后被过滤或去重的 action 上;
- 每行还携带
notification_uuid(workflow_id → notification_uuid 映射),用于跨任务传播通知上下文。
action 任务(trigger_action)的流程:
- 要求恰好一个 event ID 或 activity ID;
- 重建
WorkflowEventData; - 加载 action 与 preferred detector;
- 创建
ActionInvocation; - 调用注册在该 action 类型上的 handler。
因此外部集成调用发生在工作流求值任务之外;action handler 必须自行处理集成被删除等预期的外部状态变化。
七、处理器参考(Processor Reference)
| Processor | 输入 | 输出 | 副作用 | 异步边界 |
|---|---|---|---|---|
process_data_packet | Packet 与 source type | Detector 评估结果 | Source 缓存读取、detector 状态、Kafka 载荷 | 由产品 producer 调用 |
process_data_source | Packet 与 source type | Packet 与 detectors | Detector-source 缓存读取 | 无 |
process_detectors | Packet 与 detector 列表 | Detector 评估结果 | PostgreSQL/Redis 状态与 Issue Platform Kafka | 求值后写 Kafka |
process_data_condition_group | 条件组与事件/值 | 评估与慢条件 | 指标与日志 | 无 |
process_workflows | WorkflowEventData | 工作流评估结果 | 缓冲写入、fire history、状态更新、action 任务 | 异步 action 任务 |
process_buffered_workflows | 延迟客户端 | 无 | Cohort 元数据与延迟任务 | 异步延迟任务 |
process_delayed_workflows | Project 与可选 batch | 无 | Snuba 查询、action 状态/历史、缓冲清理 | Snuba 与 action 任务 |
八、失败与一致性模型(Failure and Consistency Model)
- 数据库图(graph)变更与缓存失效通过
transaction.on_commit()协调,避免在重填不安全时产生陈旧缓存; - Detector 状态与 Issue Platform 发布不是原子的:PostgreSQL/Redis 状态在
process_detectors发布之前提交,因此在两者之间发生故障时,重试可能被已提交的 dedupe watermark 跳过(提交顺序见 conditions.md); - Issue Platform 发布是异步的,不存在横跨 detector 状态、Kafka、group 创建、工作流求值与外部 action 的事务;
- 工作流的频率与去重减少了重复副作用,但不对外部 provider 提供分布式 exactly-once 保证;
- 延迟项在处理成功后才被移除,因此任务重试可以恢复缓冲中的工作;
- 排队中的工作所引用的模型可能在任务执行前被删除:任务与 handler 必须把"预期中的记录缺失"当作正常的生命周期竞态处理。
九、缓存(Caches)
| 缓存 | 用途 | 近似 TTL |
|---|---|---|
| caches/detector.py | 与 data source 相连的 detectors | 20 分钟(CACHE_TTL = 60 * 20) |
| caches/workflow.py | 与 detector/environment 相连的工作流 | 1 分钟 |
| caches/action_filters.py | 工作流的 IF 组与 actions | 5 分钟 |
Detector中的默认 detector 缓存 | 默认 project/type 查找 | 10 分钟(Detector.CACHE_TTL) |
把 receiver 视为模型变更契约的一部分:新增的、会影响处理的关系通常需要显式的失效覆盖。receiver 只在提交后使受支持的"当前关系 key"失效;如果原地更新DataSource的 type/source ID 或 Workflow 的 environment,不会保留并失效所有旧缓存 key,陈旧条目可能存活到 TTL 到期。
十、可观测性与调试(Observability and Debugging)
引擎会为 detector 查找/求值、条件求值、工作流处理、延迟调度、缓存与 action 发射指标。utils/log_context.py中的日志辅助函数会附加 workflow 与 event 上下文;process_workflows全程使用log_context传递detector_id、group_id、workflow_id、action_condition_group_id等上下文。
GroupedWorkflowEvaluationResult可以产出工作流评估的结构化快照,采样与直接日志由运行时 option 与 feature flag(如organizations:workflow-engine-process-workflows-logs)控制。
排查"action 未触发"时,按顺序检查以下边界:
- source 是否映射到了启用(enabled)的 detector?
- detector 状态是否发生迁移并发布了 Issue Platform 消息?
- 产生的 issue 是否关联到预期 detector?
- post-processing 或 activity 处理是否调度了工作流求值?
- 工作流是否已连接、启用,且对当前环境有效?
- WHEN 与至少一个 IF 条件组是否通过?
- 慢条件是否被缓冲并在之后处理?
- action 是否被频率或事件级去重抑制?
- fire history 是否创建、action 任务是否被调度?
- action handler 是否因外部配置缺失/无效而拒绝执行?
十一、推荐测试(Recommended Tests)
- tests/sentry/workflow_engine/test_integration.py:覆盖 detector 与 workflow 的完整路径;
- tests/sentry/workflow_engine/handlers/detector/test_stateful.py:覆盖状态迁移与 Redis 行为;
- tests/sentry/workflow_engine/processors/test_detector.py:覆盖 detector 输出与事件 detector 选择;
- tests/sentry/workflow_engine/processors/test_workflow.py:覆盖工作流求值;
- tests/sentry/workflow_engine/processors/test_data_condition_group.py:覆盖条件逻辑与 fast/slow 拆分;
- tests/sentry/workflow_engine/processors/test_delayed_workflow.py:覆盖延迟查询与 action 触发;
- tests/sentry/workflow_engine/processors/test_schedule.py:覆盖 cohort 调度与分批;
- tests/sentry/workflow_engine/tasks/test_actions.py:覆盖 action 任务重建与分发。
总结
Sentry Workflow Engine 通过"检测流水线 → Kafka/Issue Platform → 工作流流水线 → 异步 Action"的架构,把事件检测、条件求值与外部动作彻底解耦。理解其核心在于把握住几个关键心智模型:DataPacket的类型化载荷与DetectorGroupKey分组、StatefulDetectorHandler的 dedupe watermark + 优先级阈值计数器语义、WorkflowEventData在事件/活动两条入口下的差异、延迟路径的"缓冲—cohort 调度—批量 Snuba 求值"三步曲,以及频率与 dedup key 两道动作控制。结合本文引用的源码路径与测试用例,即可在排查、扩展与新 detector 开发时准确判断数据此刻正停在哪一个异步边界上。
【免费下载链接】sentryDeveloper-first error tracking and performance monitoring项目地址: https://gitcode.com/GitHub_Trending/sen/sentry
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考