第十二章:事件摄取管道深度解析
- 12.1 摄取管道总览
- 12.2 Ingest Consumer(消费者)
- 12.3 EventManager 详解
- 12.4 事件存储层
- 12.5 分组引擎(Grouping)
- 12.6 Post-Process 流水线
- 12.7 流入过滤器(Inbound Filters)
- 12.8 速率限制(Rate Limiting)
- 12.9 数据清理与脱敏(Data Scrubbing)
- 12.10 摄取性能优化
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_chunk、attachment、event、user_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_event → process_event → save_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.py 的 decode_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 详解
EventManager(src/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 端只是一个薄包装。标准化的内容包括:
- 协议版本适配:将旧版本协议字段映射到当前版本。
- 事件时间戳校验:确保
timestamp和received在合理范围内。 - 用户代理解析:当
normalize_user_agent=True时,将 User-Agent 字符串解析为浏览器、操作系统、设备信息。 - 堆栈跟踪标准化:对
stacktrace、exception等接口进行规范化。 - 标签值裁剪:裁剪过长的标签值。
- 分组配置注入:将 grouping_config 写入事件数据,供后续分组引擎使用。
- 数据类型检查:验证事件类型(
type字段),但保留generic和feedback类型不变。
一个重要细节:标准化后,generic 和 feedback 类型的事件会被恢复(因为 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),完整流程如下:
第一阶段:数据抽取与标签处理
_get_or_create_release_many:提取事件的release字段,创建或查找对应的 Release 记录。_get_event_user_many:提取用户信息(id、email、username、ip_address)。_derive_tags_many:运行自动标签推导器(例如sentry:user、sentry:release、sentry:dist等)。_derive_interface_tags_many:从接口数据中提取标签(例如从 HTTP 接口提取url、method)。_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_many 在 event_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 # 通用事件
事件分发过程:
_eventstream_insert_many调用eventstream.backend.insert()KafkaEventStream.insert()确定目标 Topic 并调用_send()_send()使用 Kafka Producer 发送消息,消息键为project_id(保证同项目事件有序),消息体为 JSON 序列化的元组(protocol_version, operation, extra_data)- 非错误事件(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编码的增强规则
配置加载通过 PrimaryGroupingConfigLoader(api.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_grouping(hashing.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_group(event_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_grouphashes(event_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 的主要职责:
- Group 状态更新:更新 Group 的
times_seen、last_seen、first_seen,管理GroupStatus状态转换(RESOLVED -> UNRESOLVED 自动回归)。 - GroupInbox 管理:控制 Issue 是否出现在 Inbox 中。
- GroupOwner 分配:运行 Issue Owners 规则,自动分配负责人。
- GroupHistory 记录:记录状态变更历史。
- GroupSnooze 检查:处理已暂缓的 Issue。
- Service Hooks:触发 Webhook 通知(
error.created等事件)。 - 插件回调:运行已安装集成插件的回调。
- SDK Crash Detection:检测 SDK 自身的崩溃。
- Replay 事件关联:将错误事件与 Session Replay 关联。
- 邮件/通知:发送 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),包括 ie、edge、safari、firefox、chrome、opera、android、opera_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 中。脱敏在两个阶段执行:
- Relay 端(事件接收时):第一轮脱敏,在数据进入管道之前。
- 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_data(datascrubbing.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_settings(datascrubbing.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 记录来减少网络往返次数,最大化吞吐量。