ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Sentry Workflow Engine 执行链路深度解析:从 DataPacket 到 Action 的完整数据流

Sentry Workflow Engine 执行链路深度解析:从 DataPacket 到 Action 的完整数据流 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 PlatformKafka 主题由 Issue Platform 消化生成的 issue 事件随后通过post-processing或issue activity进入工作流处理。这种解耦带来的直接结果是Kafka 边界意味着检测评估成功并不等价于group 创建完成、工作流已处理、Action 已派发。检测与动作执行之间存在天然异步这也是本文后面失败与一致性模型一节讨论的主题。另外普通错误事件可以不经过任何 Workflow Engine 检测器直接在 post-processing 阶段进入工作流流水线——这正是错误工作流error workflows与 issue-stream 工作流能共用同一套工作流管线的原因。相关的持久化关系请参阅>DataPacket(source_idstr(source.id), packetproduct_payload)要点packet载荷保持产品特定product-specificWorkflow Engine 不解析其内部结构生产方同时提供source typequery_type用于解析source_id到底对应哪个DataSource。2. 解析 Detectorprocess_data_packet是通用的入口函数它把两个预处理步骤串联起来并定义了可用于监控引擎健康度的指标漏斗workflow_engine.process_data_sources→workflow_engine.process_detectors→workflow_engine.process_workflowsprocess_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 enabledTrue through the default manager源码细节bulk_fetch_enabled_detectors先通过get_detectors_by_data_source走缓存再在缓存层之外过滤d.enabledreturn [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. 逐个评估 Detectorprocess_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 可能产生没有评估结果没有条件被触发一个结果未分组 detectorDetectorGroupKey为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/Ohandler 对分组值grouped values执行批量状态读写。具体的事务与 Redis pipeline 行为由同模块的DetectorStateManager实现关键点如下get_state_data先批量拉取DetectorStatePostgreSQL行再通过单个 Redis pipeline一次性GET所有 dedupe 与 counter key见 stateful.py 中的get_redis_keys_for_group_keys/bulk_get_redis_values提交时同样用 pipeline 批量SET带exREDIS_TTL或DELETE对DetectorState则用bulk_create/bulk_update_bulk_commit_redis_state/_bulk_commit_detector_state。关键语义务必理解Dedupe 值必须是正整数且按事件顺序递增。只要 Redis watermark 存在等于或更小的值都会被忽略watermark 缺失时默认 07 天后过期源码中REDIS_TTL int(timedelta(days7).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 Platformprocess_detectors将结果交给create_issue_platform_payload与produce_occurrence_to_kafka载荷类型PayloadType.OCCURRENCE/PayloadType.STATUS_CHANGE区分 occurrence 与状态变更Issue Platform 消化后创建或更新 groupoccurrence evidence 中的detector ID允许消化逻辑创建DetectorGroup关联恢复resolution使用稳定的 issue fingerprint找到同一个 group。相关实现还有associate_new_group_with_detector/ensure_association_with_detector见 processors/detector.py当 detector 已被删除时会创建detector_idNone的DetectorGroup占位以显式表达曾关联到一个已不存在的 detector。三、工作流入向Workflow Ingress事件后处理Event post-processingprocess_workflow_engine只对非 reprocessed 的 job、且事件的 group 当前处于未解决unresolved状态时调度process_workflows_event任务。任务通过 tasks/utils.py 从 EventStore 与 Issue Platform 数据重建WorkflowEventData。基于 resolution 的处理则改由受支持的 issue activity 进入。Issue activitiesIssue 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对应的 issueGroup可选的 group 状态与升级escalation信息可选的工作流环境workflow environment一个每事件本地缓存per-event local cache用于重复查询去重。但两条路径调度并不完全相同活动处理没有工作流环境因此只考虑全局globalenvironmentNone工作流活动不进入批量的条件数据获取一个 activity 的 WHEN 或 IF 结果如果还需要慢查询slow query不会继续进入延迟处理源码中evaluate_workflow_triggers与evaluate_workflows_action_filters均对Activity分支直接metrics_incr(process_workflows.enqueue_workflow.activity)并跳过入队。四、工作流流水线Workflow Pipelineprocess_workflows是工作流任务内部的主同步编排器process_workflows还会以WorkflowEvaluationOutcomeNO_DETECTOR/ENVIRONMENT_NOT_FOUND/NO_WORKFLOWS/COMPLETED返回结果方便上层统计与观测。1. 解析事件 Detectorget_detectors_for_event_data构建EventDetectors对象对于带 occurrence 的GroupEventevent_detector只从 occurrence 的detector_idevidence 解析issue_occurrence.evidence_data.get(detector_id)ID 缺失或被删除 → 没有事件 detector对于没有 occurrence 的普通错误事件使用项目默认错误 detectorDetector.get_error_detector_for_projectissue_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__抛ValueErrorget_detectors_for_event_data捕获后返回None见 processors/detector.py。preferred detector 之后会提供给需要 detector 上下文的 action。由于 detector 删除与较老的 group 可能遗留缺失的关联所有查找路径都被设计为可安全降级degrade safely。2. 解析环境与工作流处理器先解析事件环境再由get_workflows_by_detectors加载与任一选中 detector 相连的工作流。合格的工作流必须enabledTrue是全局的environmentNone或连接到事件的 environment。缓存按detector environment键控关系 receiver 在提交committed后使其失效。源码中还保留了一条 feature-flag 关闭缓存时的直接 DB 查询回退路径_get_associated_workflowsWorkflow.objects.filter(environment_filter, detectorworkflow__detector_id__in..., enabledTrue)。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 批量执行昂贵的条件缓冲BufferingDelayedWorkflowClient将 project ID 存入分片的有序集合sorted set每个 project 的工作存入一个 hashhash 字段以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 8project 被随机放入某个分片rollout 由 optiondelayed_workflow.rollout控制。调度Schedulingschedule_delayed_workflows任务在全局调度锁内执行process_buffered_workflows通过ProjectChooser把 project 划分到 cohorthashlib.sha256(project_id).digest()[0] % num_cohorts按must-process超过 1 分钟→ 最陈旧 cohort的策略选出待处理 projecttarget_max_age 1 minutemin_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按小于等于该时间戳条件删除已处理成员。评估Evaluationprocess_delayed_workflows的执行步骤解析缓冲项EventRedisData容忍坏条目时continue_on_errorTrue丢弃对已删除 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_groupsfire history 标记is_delayedTrue删除成功处理的 Redis 项cleanup_redis_buffer。容错细节Snuba 限流RateLimitExceeded由任务重试最后一次尝试仍遇限流时记录日志并按 false 处理避免整个任务失败某个 group 缺少查询结果 → 产生tainted 的非触发评估不 firing慢条件没有注册查询 handler时不生成查询处理可以正常完成并移除缓冲项不构造 tainted 结果格式损坏的 hash 条目会被排除在解析出的事件 key 集合之外不会被正常成功清理删除会一直保留到过期或人工清理。六、Action 选择与分发Action Selection and Dispatchprocessors/action.py施加两道不同的控制重复频率Repeat frequencyWorkflowActionGroupStatus记录每个(workflow, action, issue group)元组最近一次通过频率检查的时间。它在事件级去重与任务执行之前被更新频率间隔来自Workflow.config.get(frequency, 0)单位分钟见process_workflow_action_group_statuses中的frequency * timedelta(minutes1)当now - date_updated未超过频率间隔时该元组被抑制suppressed——即使之后有等价的 action 被去重或外部 handler 执行失败该元组仍被抑制缺失状态missing status会以ON CONFLICT (workflow_id, action_id, group_id) DO NOTHING批量插入插入遇IntegrityError时切换为带 EXISTS 检查的 SQL 重试以跳过引用已删除的孤儿行。事件级去重Event-level deduplicationget_unique_active_actions两步走先剔除status ! ACTIVE的 action再对每个Action.get_dedup_key()只保留一个具体 action。这是一个实现级 key并不保证所有语义等价的外部副作用都能被识别。调用的归属attribution来自该 key 保留的那个具体 action不暗示工作流执行顺序。历史与任务分发History and task dispatch处理器在fire_actions之前更新频率状态并创建WorkflowFireHistory行然后调度trigger_actionfire-history 行只描述已调度不代表外部投递成功——它们可以存在于之后被过滤或去重的 action 上每行还携带notification_uuidworkflow_id → notification_uuid 映射用于跨任务传播通知上下文。action 任务trigger_action的流程要求恰好一个 event ID 或 activity ID重建WorkflowEventData加载 action 与 preferred detector创建ActionInvocation调用注册在该 action 类型上的 handler。因此外部集成调用发生在工作流求值任务之外action handler 必须自行处理集成被删除等预期的外部状态变化。七、处理器参考Processor ReferenceProcessor输入输出副作用异步边界process_data_packetPacket 与 source typeDetector 评估结果Source 缓存读取、detector 状态、Kafka 载荷由产品 producer 调用process_data_sourcePacket 与 source typePacket 与 detectorsDetector-source 缓存读取无process_detectorsPacket 与 detector 列表Detector 评估结果PostgreSQL/Redis 状态与 Issue Platform Kafka求值后写 Kafkaprocess_data_condition_group条件组与事件/值评估与慢条件指标与日志无process_workflowsWorkflowEventData工作流评估结果缓冲写入、fire history、状态更新、action 任务异步 action 任务process_buffered_workflows延迟客户端无Cohort 元数据与延迟任务异步延迟任务process_delayed_workflowsProject 与可选 batch无Snuba 查询、action 状态/历史、缓冲清理Snuba 与 action 任务八、失败与一致性模型Failure and Consistency Model数据库图graph变更与缓存失效通过transaction.on_commit()协调避免在重填不安全时产生陈旧缓存Detector 状态与 Issue Platform 发布不是原子的PostgreSQL/Redis 状态在process_detectors发布之前提交因此在两者之间发生故障时重试可能被已提交的 dedupe watermark 跳过提交顺序见 conditions.mdIssue Platform 发布是异步的不存在横跨 detector 状态、Kafka、group 创建、工作流求值与外部 action 的事务工作流的频率与去重减少了重复副作用但不对外部 provider 提供分布式 exactly-once 保证延迟项在处理成功后才被移除因此任务重试可以恢复缓冲中的工作排队中的工作所引用的模型可能在任务执行前被删除任务与 handler 必须把预期中的记录缺失当作正常的生命周期竞态处理。九、缓存Caches缓存用途近似 TTLcaches/detector.py与 data source 相连的 detectors20 分钟CACHE_TTL 60 * 20caches/workflow.py与 detector/environment 相连的工作流1 分钟caches/action_filters.py工作流的 IF 组与 actions5 分钟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的 detectordetector 状态是否发生迁移并发布了 Issue Platform 消息产生的 issue 是否关联到预期 detectorpost-processing 或 activity 处理是否调度了工作流求值工作流是否已连接、启用且对当前环境有效WHEN 与至少一个 IF 条件组是否通过慢条件是否被缓冲并在之后处理action 是否被频率或事件级去重抑制fire history 是否创建、action 任务是否被调度action handler 是否因外部配置缺失/无效而拒绝执行十一、推荐测试Recommended Teststests/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),仅供参考
返回列表