znlgis 博客

GIS开发与技术分享 — GDAL · GeoServer · PostGIS · QGIS · OpenLayers · Cesium · FreeCAD · NPOI

第十二章:事件摄取管道深度解析


12.1 摄取管道总览

Sentry 的事件摄取管道是一条高性能、多阶段的流水线,负责将 SDK 上报的原始事件数据经过验证、标准化、脱敏、分组、存储等环节,最终落入 Snuba(ClickHouse)供查询。整个管道从 Relay 接收事件开始,经过 Kafka 消息队列、Ingest Consumer 消费、EventManager 处理,最终通过 Nodestore 和 Eventstream 完成持久化。

以下是完整的数据流图(从 SDK 到 Snuba):

SDK 上报
  │
  ▼
Relay (数据验证、PII脱敏、Rate Limiting、过滤器执行)
  │
  ├──────────► Kafka Topic: ingest-events (普通事件)
  ├──────────► Kafka Topic: ingest-attachments (附件事件)
  ├──────────► Kafka Topic: ingest-transactions (事务事件)
  ├──────────► Kafka Topic: ingest-feedback-events (反馈事件)
  │
  ▼
Ingest Consumer (Arroyo 消费者框架)
  │
  ├─ msgpack 解码
  ├─ 去重检查 (Memcached, TTL=1h)
  ├─ Killswitch 减载
  ├─ 附件存储 (Attachment Cache)
  ├─ 事件存入 Event Processing Store (Redis)
  │
  ▼
preprocess_event (异步 Celery 任务)
  │
  ├─ find_stacktraces_in_data (查找堆栈)
  ├─ submit_symbolicate → Symbolicator (符号化)
  │   │
  │   └──► do_process_event (返回后)
  │         ├─ 第二轮 Data Scrubbing
  │         ├─ Event Preprocessors (事件预处理器)
  │         └─ normalize_event (重新标准化)
  │
  ▼
save_event (异步 Celery 任务)
  │
  ├─ EventManager.save()
  │   ├─ normalize (StoreNormalizer Rust 实现)
  │   ├─ _pull_out_data (提取标签、创建Event实例)
  │   ├─ _get_or_create_release_many (发布版本)
  │   ├─ _derive_tags_many (自动标签推导)
  │   ├─ assign_event_to_group (分组引擎)
  │   │   ├─ run_primary_grouping (主分组哈希计算)
  │   │   ├─ maybe_run_secondary_grouping (辅助分组)
  │   │   ├─ maybe_check_seer_for_matching_grouphash (Seer智能匹配)
  │   │   └─ create_group_with_grouphashes (创建新Group)
  │   ├─ _nodestore_save_many (存入Nodestore)
  │   └─ _eventstream_insert_many (写入Eventstream)
  │
  ▼
Eventstream → Kafka Topics: events / transactions / eventstream-generic
  │
  ▼
Snuba Consumer (ClickHouse 写入)
  │
  ▼
post_process_group (异步 Celery 任务)
  │
  ├─ Service Hooks (Webhook 通知)
  ├─ 插件回调
  ├─ Issue Owners 分配
  ├─ 邮件通知 / Slack 通知
  ├─ SDK Crash Detection
  └─ Replay 事件关联

管道有四个核心的 Kafka Topic:

  • ingest-events:接收普通错误事件(Errors)。
  • ingest-attachments:接收带附件的事件(minidump 等),在内部又细分为 attachment_chunkattachmenteventuser_report 四种消息类型。
  • ingest-transactions:接收事务/性能事件(Transactions)。
  • ingest-feedback-events:接收用户反馈事件(Feedback)。

12.2 Ingest Consumer(消费者)

12.2.1 Kafka Topic 架构

Sentry 使用 Arroyo(Sentry 自研的 Kafka 消费者框架)来消费上述四个 Topic。每个 Topic 有独立的消费者策略工厂,定义在 src/sentry/ingest/consumer/factory.py 中:

# factory.py:78-97
class IngestStrategyFactory(ProcessingStrategyFactory[KafkaPayload]):
    def __init__(
        self,
        consumer_type: str,
        reprocess_only_stuck_events: bool,
        reprocess_only_events_not_in_nodestore: bool,
        stop_at_timestamp: int | None,
        num_processes: int,
        max_batch_size: int,
        max_batch_time: int,
        input_block_size: int | None,
        output_block_size: int | None,
    ):

此外还有专门的 IngestTransactionsStrategyFactory(位于同一文件的第198行),用于事务事件。

12.2.2 消费者类型与分工

所有消费者类型统一定义在 src/sentry/ingest/types.py 中:

# types.py:4-12
class ConsumerType(StrEnum):
    Events = "events"          # 普通错误事件 (ingest-events)
    Attachments = "attachments" # 附件事件 (ingest-attachments)
    Transactions = "transactions" # 事务事件 (ingest-transactions)
    Feedback = "feedback"       # 用户反馈 (ingest-feedback-events)

每种消费者类型的处理路径不同:

消费者类型 事件类型 处理路径
Events error decode → 去重 → Redis存储 → preprocess_eventprocess_eventsave_event
Transactions transaction decode → 去重 → Redis存储 → 直接 save_event_transaction(跳过预处理器)
Attachments attachment_chunk + event 两阶段:先处理 chunk,再处理 event(保证依赖顺序)
Feedback feedback decode → 直接 save_event_feedback(数据直传,不经过 Redis)

12.2.3 消息格式与解码

ingest-events Topic 为例,Kafka 消息体使用 msgpack 编码,解码后的消息结构如下(参见 src/sentry/ingest/consumer/simple_event.py:26-44):

# simple_event.py:54-55
message: IngestMessage = msgpack.unpackb(raw_payload, use_list=False)
# 消息结构:
# {
#     "type": "event",
#     "project_id": 123,
#     "event_id": "abc...",
#     "payload": b'{"event_id":"abc...","project":123,...}',  # JSON 字节
#     "start_time": 1691234567.890,
#     "remote_addr": "1.2.3.4",
#     "attachments": [...]  # 可选
# }

事件的实际 JSON 体在 payload 字段中,以原始字节存储,避免二次序列化。

对于 ingest-attachments Topic,消息类型更丰富,在 attachment_event.pydecode_and_process_chunks 函数中分类型处理:

# attachment_event.py:47-55
if message["type"] == "attachment_chunk":
    # 附件分片,先于 event 处理
    if not reprocess_only_stuck_events:
        process_attachment_chunk(message)
    return None  # 已消费,不进入第二步
# 返回其他类型消息进入第二步处理

12.2.4 消费者工厂与处理策略

IngestStrategyFactory.create_with_partitions() 方法构建了完整的 Arroyo 处理策略链:

# factory.py:118-190
def create_with_partitions(self, commit, partitions):
    final_step = CommitOffsets(commit)  # 提交 Kafka offset
    
    if not self.is_attachment_topic:
        # 普通 Topic: 单阶段流水线
        event_function = partial(
            process_simple_event_message,
            consumer_type=self.consumer_type, ...
        )
        next_step = maybe_multiprocess_step(mp, event_function, final_step, self._pool)
        return maybe_backpressure_step(
            health_checker=self.health_checker,
            next_step=next_step,
            is_recovery=self.reprocess_only_stuck_events,
        )
    
    # Attachments Topic: 两阶段流水线
    # Step 1: 处理 attachment_chunk(必须先完成)
    # Step 2: 处理 event / attachment / user_report
    step_2 = maybe_multiprocess_step(mp, processing_function, final_step, self._attachments_pool)
    filter_step = FilterStep(function=filter_fn, next_step=step_2)
    step_1 = maybe_multiprocess_step(mp, attachment_function, filter_step, self._pool)
    return maybe_backpressure_step(...)

两阶段设计的关键原因在于 attachment_chunk 必须在依赖它的 event 之前被处理,否则事件在尝试组装附件时会找不到分片数据。FilterStep 用于过滤掉已经处理掉的 None 消息(chunk 在处理后返回 None),避免空消息进入第二阶段。

12.2.5 去重保护

process_event 函数中(src/sentry/ingest/consumer/processors.py:117-131),每个事件在开始处理前会执行去重检查:

# processors.py:117-131
with start_span(op="deduplication_check", name="deduplication_check"):
    deduplication_key = f"ev:{project_id}:{event_id}"
    cached_value = cache.get(deduplication_key)
    if cached_value is not None:
        logger.warning(
            "pre-process-forwarder detected a duplicated event with id:%s for project:%s.",
            event_id, project_id,
        )
        return  # 已处理过,跳过

处理完成后设置1小时的缓存:

# processors.py:309-310
cache.set(deduplication_key, "", CACHE_TIMEOUT)  # CACHE_TIMEOUT=3600

注意代码中的注释明确指出:该机制使用 Memcached,无一致性保证,1小时的 TTL 实际上无法提供可靠的去重保障。它主要防止消费者重启循环引起的重复处理。

12.2.6 Killswitch 减载机制

消费者在多个关键节点检查 Killswitch,这是一种动态减载开关。当 Sentry 运维团队需要临时关闭某些项目的处理以应对紧急情况时使用:

# processors.py:133-143
if killswitch_matches_context(
    "store.load-shed-pipeline-projects",
    {
        "project_id": project_id,
        "event_id": event_id,
        "has_attachments": bool(attachments),
    },
):
    return  # 静默丢弃

还有第二个 Killswitch 在 JSON 解析后执行(processors.py:165-175),支持按 organization_id、event_type 等更细粒度的减载。


12.3 EventManager 详解

EventManagersrc/sentry/event_manager.py:343-631)是整个摄取管道的核心类,负责事件的标准化(normalization)和存储(save)。

12.3.1 EventManager 初始化

# event_manager.py:343-384
class EventManager:
    def __init__(
        self,
        data: MutableMapping[str, Any],
        version: str = "5",
        project: Project | None = None,
        grouping_config: GroupingConfig | None = None,
        client_ip: str | None = None,
        user_agent: str | None = None,
        auth: Any | None = None,
        key: Any | None = None,
        content_encoding: str | None = None,
        is_renormalize: bool = False,
        remove_other: bool | None = None,
        project_config: Any | None = None,
        sent_at: datetime | None = None,
    ):

关键参数说明:

参数 说明
data 事件原始数据字典
version 协议版本,默认为 "5"
grouping_config 分组配置,优先级:显式传入 > project_config > project 查询
is_renormalize 是否为重新标准化(重处理场景)
project_config Relay 下发的项目配置,包含 filterSettings 等
sent_at Relay 转发时间戳

12.3.2 normalize() 标准化流程

# event_manager.py:386-425
def normalize(self, project_id: int | None = None) -> None:
    with metrics.timer("events.store.normalize.duration"):
        self._normalize_impl(project_id=project_id)

def _normalize_impl(self, project_id: int | None = None) -> None:
    # ...
    from sentry_relay.processing import StoreNormalizer

    rust_normalizer = StoreNormalizer(
        project_id=self._project.id if self._project else project_id,
        client_ip=self._client_ip,
        client=self._auth.client if self._auth else None,
        key_id=str(self._key.id) if self._key else None,
        grouping_config=self._grouping_config,
        protocol_version=str(self.version) if self.version is not None else None,
        is_renormalize=self._is_renormalize,
        remove_other=self._remove_other,
        normalize_user_agent=True,
        sent_at=self.sent_at.isoformat() if self.sent_at is not None else None,
        json_dumps=orjson.dumps,
        **DEFAULT_STORE_NORMALIZER_ARGS,
    )

    self._data = rust_normalizer.normalize_event(dict(self._data), json_loads=orjson.loads)

标准化过程的核心实现在 Rust 库 sentry_relay.processing.StoreNormalizer 中,Python 端只是一个薄包装。标准化的内容包括:

  1. 协议版本适配:将旧版本协议字段映射到当前版本。
  2. 事件时间戳校验:确保 timestampreceived 在合理范围内。
  3. 用户代理解析:当 normalize_user_agent=True 时,将 User-Agent 字符串解析为浏览器、操作系统、设备信息。
  4. 堆栈跟踪标准化:对 stacktraceexception 等接口进行规范化。
  5. 标签值裁剪:裁剪过长的标签值。
  6. 分组配置注入:将 grouping_config 写入事件数据,供后续分组引擎使用。
  7. 数据类型检查:验证事件类型(type 字段),但保留 genericfeedback 类型不变。

一个重要细节:标准化后,genericfeedback 类型的事件会被恢复(因为 Rust Normalizer 不认识这些类型):

# event_manager.py:421-425
if pre_normalize_type in ("generic", "feedback"):
    self._data["type"] = pre_normalize_type

12.3.3 save() 存储总流程

# event_manager.py:431-508
def save(
    self,
    project_id: int | None = None,
    project: Project | None = None,
    raw: bool = False,
    assume_normalized: bool = False,
    start_time: float | None = None,
    cache_key: str | None = None,
    skip_send_first_transaction: bool = False,
    attachments: list[CachedAttachment] | None = None,
) -> Event:

根据事件类型分流处理:

    event_type = self._data.get("type")
    if event_type == "transaction":
        # 事务事件:走 save_transaction_events 路径
        jobs = save_transaction_events([job], projects, skip_send_first_transaction)
        return jobs[0]["event"]
    elif event_type == "generic":
        # 通用事件(Issue Platform):走 save_generic_events 路径
        jobs = save_generic_events([job], projects)
        return jobs[0]["event"]
    else:
        # 错误事件:走 save_error_events 路径
        return self.save_error_events(project, job, projects, metric_tags, attachments or [], raw, cache_key)

12.3.4 save_error_events() 错误事件保存

这是最复杂的保存路径(event_manager.py:510-631),完整流程如下:

第一阶段:数据抽取与标签处理

  1. _get_or_create_release_many:提取事件的 release 字段,创建或查找对应的 Release 记录。
  2. _get_event_user_many:提取用户信息(id、email、username、ip_address)。
  3. _derive_tags_many:运行自动标签推导器(例如 sentry:usersentry:releasesentry:dist 等)。
  4. _derive_interface_tags_many:从接口数据中提取标签(例如从 HTTP 接口提取 urlmethod)。
  5. _derive_client_error_sampling_rate:提取客户端的错误采样率。

第二阶段:分组引擎

assign_event_to_group 是核心函数(后文 12.5 节详述),它会:

  • 计算事件的主分组哈希(Primary Grouping)
  • 尝试匹配已有的 GroupHash
  • 如未匹配,尝试辅助分组配置(Secondary Grouping)
  • 仍未匹配则调用 Seer ML 服务
  • 最终创建新 Group 或更新已有 Group

第三阶段:环境与发布关联

  • _get_or_create_environment_many:创建/查找 Environment
  • _get_or_create_group_environment_many:关联 Group 与 Environment
  • _get_or_create_release_associated_models:创建 ReleaseEnvironment、ReleaseProjectEnvironment
  • _increment_release_associated_counts_many:递增新分组计数
  • _get_or_create_group_release_many:创建 GroupRelease 关联

第四阶段:持久化

  • _tsdb_record_all_metrics:向 TSDB(时序数据库)写入指标计数
  • _materialize_event_metrics:具体化事件指标
  • _nodestore_save_many:将事件体写入 Nodestore
  • _eventstream_insert_many:将事件元数据写入 Kafka Eventstream,触发 Snuba 写入和 Post-Process

第五阶段:后处理信号

  • first_event_received 信号:项目首个事件
  • first_event_with_minified_stack_trace_received:首个含压缩堆栈的事件
  • 重处理事件的旧主哈希清理

12.4 事件存储层

12.4.1 Nodestore —— 事件体存储

Nodestore 负责存储事件的完整 JSON 体。它是一个抽象层,支持多种后端:

src/sentry/nodestore/
├── __init__.py       # 后端选择逻辑
├── base.py           # NodeStorage 基类
├── bigtable/         # Google Bigtable 后端 (SaaS)
├── django/           # Django ORM 后端 (开发/小规模)
└── filesystem/       # 文件系统后端 (测试)

_nodestore_save_manyevent_manager.py:1076-1104 中实现:

# event_manager.py:1076-1104
def _nodestore_save_many(jobs: Sequence[Job], app_feature: str) -> None:
    inserted_time = datetime.now(timezone.utc).timestamp()
    for job in jobs:
        subkeys = {}
        event = job["event"]
        # 错误事件:保留未处理的事件数据作为 subkey
        if event.get_event_type() not in ("transaction", "generic") and job["groups"]:
            unprocessed = event_processing_store.get(
                cache_key_for_event({"project": event.project_id, "event_id": event.event_id}),
                unprocessed=True,
            )
            if unprocessed is not None:
                subkeys["unprocessed"] = unprocessed

        job["event"].data["nodestore_insert"] = inserted_time
        job["event"].data.save(subkeys=subkeys)

关键细节:

  • 错误事件会额外存储一个 unprocessed 子键,包含重处理所需的信息。
  • 事务和通用事件不存储 unprocessed 子键。
  • nodestore_insert 时间戳记录事件写入 Nodestore 的时间。

12.4.2 Eventstream —— 事件流分发

Eventstream 是一个服务抽象层,负责将已保存的事件分发到下游消费者。src/sentry/eventstream/base.py 定义了基类 EventStream

# base.py:39-195
class EventStream(Service):
    __all__ = (
        "insert",           # 插入事件
        "start_delete_groups", "end_delete_groups",  # 删除分组
        "start_merge", "end_merge",                   # 合并分组
        "start_unmerge", "end_unmerge",               # 拆分分组
        "start_delete_tag", "end_delete_tag",         # 删除标签
        "tombstone_events_unsafe",                     # 标记事件删除
        "replace_group_unsafe",                       # 替换分组
        "exclude_groups",                              # 排除分组
        "requires_post_process_forwarder",
    )

Kafka 后端的实现在 src/sentry/eventstream/kafka/backend.py:34-235

# kafka/backend.py:34-37
class KafkaEventStream(SnubaProtocolEventStream):
    def __init__(self, **options: Any) -> None:
        self.topic = Topic.EVENTS                # 错误事件
        self.transactions_topic = Topic.TRANSACTIONS  # 事务事件
        self.issue_platform_topic = Topic.EVENTSTREAM_GENERIC  # 通用事件

事件分发过程:

  1. _eventstream_insert_many 调用 eventstream.backend.insert()
  2. KafkaEventStream.insert() 确定目标 Topic 并调用 _send()
  3. _send() 使用 Kafka Producer 发送消息,消息键为 project_id(保证同项目事件有序),消息体为 JSON 序列化的元组 (protocol_version, operation, extra_data)
  4. 非错误事件(Transaction、Generic)使用随机分区以均匀负载

Post-Process Forwarder 机制

Base 类的 insert() 方法还负责触发 post_process 任务:

# base.py:60-92
def _dispatch_post_process_group_task(self, event_id, project_id, group_id, ...):
    cache_key = cache_key_for_event({"project": project_id, "event_id": event_id})
    post_process_group.apply_async(
        kwargs={
            "is_new": is_new,
            "is_regression": is_regression,
            "is_new_group_environment": is_new_group_environment,
            "primary_hash": primary_hash,
            "cache_key": cache_key,
            "group_id": group_id,
            "group_states": group_states,
            "occurrence_id": occurrence_id,
            "project_id": project_id,
            "eventstream_type": eventstream_type,
        },
        headers={"sentry-propagate-traces": False},
    )

注意 Kafka 后端的 requires_post_process_forwarder() 返回 True,表示它需要一个独立的 Forwarder 进程来消费 Kafka 消息并触发 post_process(而非在 insert() 中直接触发)。

12.4.3 save_event 全流程串联

src/sentry/tasks/store.py 中的任务函数将整个事件处理流程串联起来:

Kafka消息到达
    │
    ▼
Ingest Consumer (simple_event.py / attachment_event.py)
    │
    ├─ 解码 & 去重
    ├─ Redis 存储 (event_processing_store / transaction_processing_store)
    │
    ▼
preprocess_event (store.py:231) [Celery任务]
    │
    ├─ 从Redis读取事件数据
    ├─ find_stacktraces_in_data (查找堆栈跟踪)
    ├─ get_symbolication_functions (确定需符号化的平台)
    ├─ submit_symbolicate → Symbolicator (原生符号化)
    │
    ▼
process_event (store.py:423) [Celery任务]
    │
    ├─ is_process_disabled (Killswitch检查)
    ├─ scrub_data (第二轮PII脱敏,因符号化可能引入敏感数据)
    ├─ get_event_preprocessors (事件预处理器)
    └─ normalize_event (重新标准化)
    │
    ▼
save_event (store.py:623) [Celery任务]
    │
    └─ _do_save_event
        ├─ EventManager(data).save()
        │   ├─ normalize (StoreNormalizer)
        │   ├─ _pull_out_data (标签提取)
        │   ├─ assign_event_to_group (分组)
        │   ├─ _nodestore_save_many (Nodestore)
        │   └─ _eventstream_insert_many (Eventstream)
        └─ 清理 (Redis删除、附件清理、重处理标记)

对于事务事件(Transactions),流程更短:

Ingest Consumer → Redis存储 → save_event_transaction (直接跳过预处理)

12.5 分组引擎(Grouping)

分组引擎是 Sentry 最核心的算法之一,它决定了哪些事件属于同一个 Issue。分组过程将相似的事件聚合成一个 Group,使得开发者能够在一个 Issue 页面看到所有相关的错误发生情况。

12.5.1 分组配置(GroupingConfig)

分组配置在 src/sentry/grouping/api.py:64-67 中定义为 TypedDict:

class GroupingConfig(TypedDict):
    id: str              # 分组策略ID,如 "newstyle:2023-01-11"
    enhancements: str    # Base64编码的增强规则

配置加载通过 PrimaryGroupingConfigLoaderapi.py:136-140):

class PrimaryGroupingConfigLoader(ProjectGroupingConfigLoader):
    option_name = "sentry:grouping_config"
    cache_prefix = "grouping-enhancements:"

项目级别的分组配置选项 sentry:grouping_config 决定了使用哪种分组策略。当前主流策略有:

  • newstyle:2023-01-11:当前最新策略
  • newstyle:2019-10-29:较早的策略(过渡期使用)

12.5.2 分组策略与变体(Variants)

分组策略定义在 src/sentry/grouping/strategies/configurations.py 中,每种策略会产生多个 Variant(变体)。Variant 类型在 src/sentry/grouping/variants.py 中定义:

Variant 类型 说明
ComponentVariant 默认变体,由堆栈跟踪的关键帧组成
CustomFingerprintVariant 用户自定义指纹规则匹配
ChecksumVariant 旧版 checksum(直接使用 checksum 字段作为哈希)
HashedChecksumVariant 新版 checksum(对 checksum 做 MD5 哈希)
SaltedComponentVariant 加盐的组件变体(用于区分不同场景)
FallbackVariant 降级变体(当其他变体都失败时使用)

每个 Variant 结构:

# variants.py:26-77
class BaseVariant(ABC):
    variant_name: str | None = None
    @property
    def contributes(self) -> bool: return True
    @property
    @abstractmethod
    def type(self) -> str: ...
    def get_hash(self) -> str | None: return None
    @property
    def key(self) -> str: return self.type
    @property
    def description(self) -> str: return self.type
    @property
    def hint(self) -> str | None: return None

12.5.3 指纹算法与自定义 Fingerprint

Sentry 的指纹规则定义在 src/sentry/grouping/fingerprinting/ 目录下:

fingerprinting/
├── __init__.py       # 入口
├── rules.py          # FingerprintRule 类
├── types.py          # FingerprintInfo 等类型
├── utils.py          # 工具函数:is_default_fingerprint_var, resolve_fingerprint_values
├── parser.py         # 规则解析器
├── matchers.py       # 匹配器
└── exceptions.py     # 异常类

默认指纹值为 ``,表示使用系统默认分组算法。用户可以在项目设置中配置自定义指纹规则,例如:

# 按异常类型分组
type:MyCustomError -> my-custom-group

# 按错误消息分组(忽略参数值)
message:"Connection to * failed" -> connection-error

服务端指纹应用在 _calculate_event_grouping 中调用(src/sentry/grouping/ingest/hashing.py:67-70):

event.data["fingerprint"] = event.data.data.get("fingerprint") or [""]
apply_server_side_fingerprinting(
    event.data.data, get_fingerprinting_config_for_project(project)
)

12.5.4 分组哈希计算流程

_calculate_event_groupinghashing.py:50-82)是哈希计算的入口:

def _calculate_event_grouping(
    project: Project, event: Event, grouping_config: GroupingConfig
) -> tuple[list[str], dict[str, BaseVariant]]:
    loaded_grouping_config = load_grouping_config(grouping_config)

    # 1. 应用服务端指纹规则
    event.data["fingerprint"] = event.data.data.get("fingerprint") or [""]
    apply_server_side_fingerprinting(
        event.data.data, get_fingerprinting_config_for_project(project)
    )

    # 2. 标准化堆栈跟踪用于分组
    event.normalize_stacktraces_for_grouping(loaded_grouping_config)

    # 3. 计算哈希和变体
    hashes, variants = event.get_hashes_and_variants(loaded_grouping_config)

    return (hashes, variants)

12.5.5 Secondary Grouping 与过渡期

当分组配置发生变更时,Sentry 使用双轨制(Secondary Grouping)来避免创建重复的 Issue。项目的过渡状态由 is_in_transition 函数检查(src/sentry/grouping/ingest/config.py)。

assign_event_to_groupevent_manager.py:1310-1374)中,如果主分组配置未匹配到现有 Group,会尝试使用辅助配置:

# 主分组
primary = get_hashes_and_grouphashes(job, run_primary_grouping, metric_tags)
if primary.existing_grouphash:
    group_info = handle_existing_grouphash(...)
else:
    # 辅助分组(过渡期使用)
    secondary = get_hashes_and_grouphashes(job, maybe_run_secondary_grouping, metric_tags)
    all_grouphashes = primary.grouphashes + secondary.grouphashes
    if secondary.existing_grouphash:
        group_info = handle_existing_grouphash(...)
    else:
        # Seer ML 匹配 or 创建新 Group

12.5.6 Seer 匹配与分组创建

如果在主分组和辅助分组中都找不到匹配,Sentry 会调用 Seer 服务(基于机器学习的相似度匹配):

# event_manager.py:1344-1353
seer_matched_grouphash = maybe_check_seer_for_matching_grouphash(
    event, primary.grouphashes[0], primary.variants, all_grouphashes
)
if seer_matched_grouphash:
    group_info = handle_existing_grouphash(job, seer_matched_grouphash, all_grouphashes)
else:
    group_info = create_group_with_grouphashes(job, all_grouphashes)

create_group_with_grouphashesevent_manager.py:1471)在数据库事务中创建新 Group:

with transaction.atomic(router.db_for_write(GroupHash)):
    # 获取锁,防止并发创建
    group = Group.objects.create(
        project=project,
        short_id=project.next_short_id(),
        **kwargs  # 包含 first_seen, data, platform, message, level, culprit 等
    )
    add_group_id_to_grouphashes(group, grouphashes)

12.6 Post-Process 流水线

12.6.1 post_process_group 任务

事件保存到 Eventstream 后,Snuba Consumer 消费 Kafka 消息写入 ClickHouse,随后触发 post_process_group Celery 任务(src/sentry/tasks/post_process.py)。该任务处理所有异步的后处理逻辑:

# post_process.py 中的核心数据结构
class PostProcessJob(TypedDict, total=False):
    event: GroupEvent
    group_state: GroupState
    is_reprocessed: bool
    has_reappeared: bool
    has_escalated: bool
    halt_post_process: bool

Post-process 的主要职责:

  1. Group 状态更新:更新 Group 的 times_seenlast_seenfirst_seen,管理 GroupStatus 状态转换(RESOLVED -> UNRESOLVED 自动回归)。
  2. GroupInbox 管理:控制 Issue 是否出现在 Inbox 中。
  3. GroupOwner 分配:运行 Issue Owners 规则,自动分配负责人。
  4. GroupHistory 记录:记录状态变更历史。
  5. GroupSnooze 检查:处理已暂缓的 Issue。
  6. Service Hooks:触发 Webhook 通知(error.created 等事件)。
  7. 插件回调:运行已安装集成插件的回调。
  8. SDK Crash Detection:检测 SDK 自身的崩溃。
  9. Replay 事件关联:将错误事件与 Session Replay 关联。
  10. 邮件/通知:发送 Issue 通知邮件。

12.6.2 回调链与插件机制

post_process 使用锁管理器防止并发处理同一个 Group:

# post_process.py:56-60
locks = LockManager(
    build_instance_from_options_of_type(
        LockBackend, settings.SENTRY_POST_PROCESS_LOCKS_BACKEND_OPTIONS
    )
)

处理过程中使用 safe_execute 包裹每个回调,确保单个回调失败不影响其他回调的执行。


12.7 流入过滤器(Inbound Filters)

流入过滤器在 Relay 端执行,用于在事件进入 Sentry 管道之前就将其丢弃。这减少了不必要的存储成本和处理开销。Sentry 的过滤器定义在 src/sentry/ingest/inbound_filters.py 中。

12.7.1 静态过滤规则

注册在 get_all_filter_specs()inbound_filters.py:63-80)中的五个内置过滤器:

过滤器 ID 名称 说明
browser-extensions 浏览器扩展过滤 过滤已知由浏览器扩展注入脚本导致的错误
localhost 本地地址过滤 过滤来自 127.0.0.1::1 的事件
legacy-browsers 旧浏览器过滤 过滤 IE11以下、Safari 15以下等旧浏览器
web-crawlers 爬虫过滤 过滤已知网络爬虫(如 Googlebot)产生的事件
filtered-transaction 健康检查事务过滤 过滤匹配常见健康检查模式的 Transaction

每个过滤器的统计指标记录在 TSDB 中:

# inbound_filters.py:37-49
FILTER_STAT_KEYS_TO_VALUES = {
    FilterStatKeys.IP_ADDRESS: TSDBModel.project_total_received_ip_address,
    FilterStatKeys.RELEASE_VERSION: TSDBModel.project_total_received_release_version,
    FilterStatKeys.ERROR_MESSAGE: TSDBModel.project_total_received_error_message,
    FilterStatKeys.BROWSER_EXTENSION: TSDBModel.project_total_received_browser_extensions,
    FilterStatKeys.LEGACY_BROWSER: TSDBModel.project_total_received_legacy_browsers,
    FilterStatKeys.LOCALHOST: TSDBModel.project_total_received_localhost,
    FilterStatKeys.WEB_CRAWLER: TSDBModel.project_total_received_web_crawlers,
    # ...
}

旧版本浏览器过滤器支持细粒度的子过滤器选择(inbound_filters.py:229-283),包括 ieedgesafarifirefoxchromeoperaandroidopera_mini 等。

12.7.2 通用过滤器(Generic Filters)

与静态过滤器不同,通用过滤器使用 Relay 的 RuleCondition DSL 表达任意过滤条件。当前激活的通用过滤器定义在 inbound_filters.py:421-425

ACTIVE_GENERIC_FILTERS: Sequence[tuple[str, Callable[[], RuleCondition | None]]] = [
    ("chunk-load-error", _chunk_load_error_filter),
    ("react-hydration-errors", _hydration_error_filter),
    ("custom-error", _custom_error_filter),
]

Chunk Load Error 过滤器inbound_filters.py:366-385):匹配 Webpack/Turbopack 的代码分片加载错误:

values = [
    ("ChunkLoadError", "Loading chunk *"),
    ("*Uncaught *", "ChunkLoadError: Loading chunk *"),
    ("ChunkLoadError", "Failed to load chunk *"),
    # ...
]

React Hydration Error 过滤器inbound_filters.py:397-417):匹配 React 服务端渲染水合错误:

values = [
    (None, "*https://reactjs.org/docs/error-decoder.html?invariant={418,419,421,422,423,425}*"),
    (None, "*https://react.dev/errors/{418,419,421,422,423,425}*"),
]

自定义错误过滤器:通过 settings.SENTRY_INBOUND_FILTER_CUSTOM_VALUES 配置自定义错误消息匹配。

12.7.3 过滤器状态管理

过滤器的开关状态存储在项目选项 filters:{filter_id} 中:

# inbound_filters.py:83-119
def set_filter_state(filter_id, project: Project, state):
    # 旧浏览器过滤器特殊处理:支持 active / subfilters
    if flt == _legacy_browsers_filter:
        # option_val 可能是 "0", "1", 或 sorted subfilters 列表
        ProjectOption.objects.set_value(
            project=project, key=f"filters:{filter_id}", value=option_val
        )
    else:
        # 所有布尔过滤器
        ProjectOption.objects.set_value(
            project=project,
            key=f"filters:{filter_id}",
            value="1" if state.get("active", False) else "0",
        )

12.8 速率限制(Rate Limiting)

12.8.1 Relay 端的速率限制

速率限制主要在 Relay 端执行,而非 Sentry 服务端。Relay 根据项目配置中的配额设置决定是否接受或拒绝事件。src/sentry/quotas/base.py 定义了配额的基础结构:

# quotas/base.py:27-34
@unique
class QuotaScope(IntEnum):
    ORGANIZATION = 1   # 组织级别
    PROJECT = 2        # 项目级别
    KEY = 3            # DSN Key 级别

@dataclass
class AbuseQuota:
    id: str
    option: str
    categories: list[DataCategory]
    scope: AbuseQuotaScope
    namespace: str | None = None

12.8.2 限流层级与配置

Sentry 支持三个层级的速度限制:

层级 范围 说明
Organization 整个组织 控制组织的总事件量
Project 单个项目 控制每个项目的配额
Key 单个 DSN Key 控制每个 DSN 的配额

每种配额都可以按 DataCategory 设置不同的限制:

DataCategory 说明
ERROR 错误事件
TRANSACTION 事务/性能事件
ATTACHMENT 附件
REPLAY Session Replay
PROFILE 性能剖析
MONITOR_SEAT 监控席位
SPAN Span 数据
TRACE_METRIC 跟踪指标
LOG_BYTE 日志字节数

配额配置会下发到 Relay,Relay 在接收到事件时进行实时限流。被限流的事件会被标记为 RATE_LIMITED Outcome 并丢弃。


12.9 数据清理与脱敏(Data Scrubbing)

数据脱敏机制定义在 src/sentry/relay/datascrubbing.py 中。脱敏在两个阶段执行:

  1. Relay 端(事件接收时):第一轮脱敏,在数据进入管道之前。
  2. Sentry 服务端(符号化后):第二轮脱敏,因为符号化过程可能引入新的敏感数据(如文件路径中的用户名)。

12.9.1 PII 配置合并

get_all_pii_configs 函数(datascrubbing.py:80-88)组合了两个配置来源:

def get_all_pii_configs(project):
    # 1. 组织和项目级别的自定义 PII 规则
    pii_config = get_pii_config(project)
    if pii_config:
        yield pii_config

    # 2. 数据脱敏设置(敏感字段、IP 脱敏等)
    settings = get_datascrubbing_settings(project)
    yield convert_datascrubbing_config(settings, json_dumps=orjson.dumps, json_loads=orjson.loads)

PII 配置合并时,组织规则优先于项目规则:

# datascrubbing.py:112-144
def _merge_pii_configs(prefixes_and_configs):
    # 为每个配置源添加前缀,确保规则名唯一
    # 组织规则前缀: "organization:"
    # 项目规则前缀: "project:"
    for prefix, partial_config in prefixes_and_configs:
        for rule_name, rule in rules.items():
            prefixed_rule_name = f"{prefix}{rule_name}"
            merged_config.setdefault("rules", {})[prefixed_rule_name] = (...)

合并策略确保组织级别的规则不被项目级别规则覆盖。

12.9.2 脱敏流程

核心函数 scrub_datadatascrubbing.py:92-109):

def scrub_data(project: Project, event: MutableMapping[str, Any]) -> MutableMapping[str, Any]:
    for config in get_all_pii_configs(project):
        event = pii_strip_event(
            config, dict(event), json_loads=orjson.loads, json_dumps=orjson.dumps
        )
    return event

# process_event 中的第二轮脱敏 (store.py:370-376):
if has_changed:
    new_data = safe_execute(scrub_data, project=project, event=data)
    if new_data is not None:
        data = new_data

实际脱敏由 Rust 库 sentry_relay.processing.pii_strip_event 执行,支持:

  • 正则匹配替换
  • 字段路径精确删除
  • 数据类型(信用卡号、邮箱、IP 地址等)检测与掩码

12.9.3 数据脱敏设置

get_datascrubbing_settingsdatascrubbing.py:49-77)读取组织和项目的脱敏设置:

def get_datascrubbing_settings(project):
    rv = {}
    # 安全字段(不被脱敏的字段)
    rv["excludeFields"] = org.get_option("sentry:safe_fields", []) + \
                           project.get_option("sentry:safe_fields", [])
    # 数据脱敏开关
    if org.get_option("sentry:require_scrub_data", False) or \
       project.get_option("sentry:scrub_data", True):
        rv["scrubData"] = True
    # IP 脱敏开关
    if org.get_option("sentry:require_scrub_ip_address", False) or \
       project.get_option("sentry:scrub_ip_address", False):
        rv["scrubIpAddresses"] = True
    # 敏感字段列表
    rv["sensitiveFields"] = org.get_option("sentry:sensitive_fields", []) + \
                             project.get_option("sentry:sensitive_fields", [])
    return rv

数据脱敏默认开启(scrubData 默认为 True),IP 脱敏默认关闭。


12.10 摄取性能优化

12.10.1 多进程处理与批处理

Arroyo 消费者框架支持多进程并行处理,通过 MultiProcessConfig 配置(factory.py:25-30):

class MultiProcessConfig(NamedTuple):
    num_processes: int          # 进程数
    max_batch_size: int         # 最大批量大小
    max_batch_time: int         # 最大批处理时间(毫秒)
    input_block_size: int | None  # 输入分块大小
    output_block_size: int | None # 输出分块大小

在策略工厂中,当 num_processes > 1 时启用多进程:

# factory.py:111-114
if num_processes > 1:
    self.multi_process = MultiProcessConfig(
        num_processes, max_batch_size, max_batch_time, input_block_size, output_block_size
    )

maybe_multiprocess_step 函数(factory.py:37-58)决定使用多进程还是单线程:

def maybe_multiprocess_step(mp, function, next_step, pool):
    if mp is not None:
        return run_task_with_multiprocessing(
            function=function, next_step=next_step,
            max_batch_size=mp.max_batch_size,
            max_batch_time=mp.max_batch_time,
            pool=pool, ...
        )
    else:
        return RunTask(function=function, next_step=next_step)

Attachments Topic 使用两个独立的进程池(factory.py:107-110),因为需要两阶段流水线。

12.10.2 背压控制(Backpressure)

Sentry 的摄取管道实现了背压机制(src/sentry/processing/backpressure/),在系统负载过高时自动减速消费者:

backpressure/
├── __init__.py
├── arroyo.py      # Arroyo集成
├── health.py      # 健康度计算(基于Redis内存使用率)
├── memory.py      # 内存监控
├── monitor.py     # 监控指标
└── topology.py    # 拓扑配置

在策略工厂中,背压步骤通过 maybe_backpressure_step 添加(factory.py:61-75):

def maybe_backpressure_step(health_checker, next_step, is_recovery):
    if is_recovery:
        return next_step  # 恢复消费者不加背压(只处理卡住的事件)
    else:
        return create_backpressure_step(health_checker=health_checker, next_step=next_step)

恢复消费者(reprocess_only_stuck_events=True)不启用背压,因为它们只处理已经在 Redis 中的”卡住”事件,不会增加 Redis 内存压力。

12.10.3 Redis 缓存与事件处理存储

事件在进入处理管道后,JSON 数据存储在 Redis 中的 Event Processing Store 中。这是一个临时存储层,避免在各阶段任务之间重复传输大型 JSON 负载:

# processors.py:154-157
if consumer_type == ConsumerType.Transactions or data.get("type") == "transaction":
    processing_store = transaction_processing_store
else:
    processing_store = event_processing_store

Task 传递时只传递 cache_key(Redis 键),而非完整事件数据:

# store.py:98-119
def submit_save_event(task_kind, project_id, cache_key, event_id, start_time, data, inline=False):
    if cache_key:
        data = None  # 不传递数据,由任务自行从Redis读取
    
    task_kwargs = {
        "cache_key": cache_key,
        "data": data,
        "start_time": start_time,
        "event_id": event_id,
        "project_id": project_id,
    }

事件处理完成后清理 Redis 中的数据:

# store.py:599-605 (transactions 立即清理)
if consumer_type == ConsumerType.Transactions and event_id:
    if cache_key:
        processing_store.delete_by_key(cache_key)

此外,Sentry 内部事件使用 cache_key_for_event 作为缓存键,用于去重和附件缓存。整个管道通过 Redis Pipeline 批量操作、批量 TSDB 记录来减少网络往返次数,最大化吞吐量。