第十三章:数据模型与存储层
- 13.1 存储层总览
- 13.2 PostgreSQL 数据模型
- 13.3 Group(Issue)模型详解
- 13.4 Nodestore
- 13.5 Snuba / ClickHouse 查询层
- 13.6 Redis / Memcached 缓存策略
- 13.7 Kafka 消息队列
- 13.8 数据库迁移管理
- 13.9 数据生命周期管理
- 13.10 存储选择指南
13.1 存储层总览
Sentry 的存储架构是一个典型的混合持久化(Polyglot Persistence)设计,通过五层存储体系实现高性能的事件采集、处理与查询:
+------------------+
| Client SDKs |
+--------+---------+
|
v
+--------+---------+
| Relay / Ingest |
+--------+---------+
|
+-----------------+-----------------+
| | |
Kafka 消息队列 PostgreSQL Nodestore
(异步解耦/缓冲) (元数据/关系) (事件体 KV)
| |
v |
+--------+---------+ |
| Snuba/ClickHouse |--------------------------+
| (时序/聚合查询) |
+--------+---------+
|
v
+------+------+
| Redis/Mcache|
| (查询缓存) |
+-------------+
各层职责明确:
| 存储层 | 技术选型 | 存储内容 | 核心特性 |
|---|---|---|---|
| PostgreSQL | Django ORM | 组织、项目、Issue 元数据、用户、团队、Release、规则配置 | 事务一致、关系模型、迁移可控 |
| Nodestore | Django PG / Bigtable | 事件原始 JSON 体(Event Blob) | KV 存储、压缩、多后端 |
| Snuba/ClickHouse | ClickHouse | 错误事件、Transactions、Spans、Profiles、Replays 等时序数据 | 列存、高压缩比、聚合查询 |
| Redis/Memcached | Redis | 查询结果缓存、计数器缓存、项目配置缓存 | 亚毫秒延迟、TTL |
| Kafka | Arroyo/Kafka | 事件流、订阅结果、摄取管道 | 异步解耦、持久化、重放 |
数据流向:SDK 事件先进入 Relay 摄取层,经过规范化(Normalization)和分组(Grouping)后,元数据写入 PostgreSQL(Group 模型),事件原始 JSON 存入 Nodestore,同时事件数据通过 EventStream 发送到 Kafka Topic,由 Snuba Consumer 消费写入 ClickHouse,最终通过 Snuba 查询层对外提供聚合查询能力。
13.2 PostgreSQL 数据模型
Sentry 的 PostgreSQL 数据库承载所有元数据和关系数据。核心模型以 Django ORM 定义,位于 src/sentry/models/ 目录下。模型通过 @cell_silo_model 装饰器标记了 Silo 归属(Control/Region)。
13.2.1 Organization 模型
Organization 是 Sentry 中的顶级命名空间,所有项目数据归属在其下。定义在 src/sentry/models/organization.py:150:
@snowflake_id_model
@cell_silo_model
class Organization(ReplicatedCellModel):
"""
An organization represents a group of individuals which maintain ownership of projects.
"""
category = OutboxCategory.ORGANIZATION_UPDATE
replication_version = 5
name = models.CharField(max_length=64)
slug = SentryOrgSlugField(unique=True)
status = BoundedPositiveIntegerField(
choices=OrganizationStatus.as_choices(),
default=OrganizationStatus.ACTIVE.value
)
date_added = models.DateTimeField(default=timezone.now)
default_role = models.CharField(max_length=32, default=str(roles.get_default().id))
is_test = models.BooleanField(default=False, db_default=False)
class flags(TypedClassBitField):
allow_joinleave: bool
enhanced_privacy: bool
disable_shared_issues: bool
early_adopter: bool
require_2fa: bool
disable_new_visibility_features: bool
# ... 更多位标志字段
设计要点:
- Snowflake ID:Organization 使用 Snowflake 算法生成全局唯一、时间有序的主键,由
@snowflake_id_model装饰器和generate_snowflake_id()函数管理。 ReplicatedCellModel:Organization 是跨 Silo 复制的模型,通过 Outbox 模式在 Control Silo 与 Region Silo 之间同步变更。- BitField flags:使用
TypedClassBitField将多个布尔标志压缩到一个整数列中,避免字段膨胀。标志位按定义顺序占位,不能插入或删除中间位(会破坏已有数据),只能追加到末尾。 OrganizationStatus定义了组织的生命周期:ACTIVE(0)->PENDING_DELETION(1)->DELETION_IN_PROGRESS(2)。
OrganizationManager 提供了多种查询方法:
class OrganizationManager(BaseManager["Organization"]):
def get_for_user_ids(self, user_ids):
"""返回一组用户有权访问的所有组织"""
return self.filter(
status=OrganizationStatus.ACTIVE,
member_set__user_id__in=user_ids,
)
def get_for_user(self, user, scope=None, only_visible=True):
"""返回单个用户有权访问的组织列表"""
# 通过 OrganizationMember 中间表关联查询
qs = OrganizationMember.objects.filter(user_id=user.id).select_related("organization")
if only_visible:
qs = qs.filter(organization__status=OrganizationStatus.ACTIVE)
# ...
13.2.2 Project 模型
Project 是 Sentry 的数据入口粒度——每个错误事件都归属到特定的 Project。定义在 src/sentry/models/project.py:234:
class Project(Model):
"""
Projects are permission based namespaces which generally
are the top level entry point for all data.
"""
slug = SentrySlugField(max_length=50)
name = models.CharField(max_length=200)
organization = FlexibleForeignKey("sentry.Organization")
teams = models.ManyToManyField("sentry.Team", related_name="teams",
through=ProjectTeam)
public = models.BooleanField(default=False)
date_added = models.DateTimeField(default=timezone.now)
status = BoundedPositiveIntegerField(
default=0,
choices=(
(ObjectStatus.ACTIVE, _("Active")),
(ObjectStatus.PENDING_DELETION, _("Pending Deletion")),
(ObjectStatus.DELETION_IN_PROGRESS, _("Deletion in Progress")),
),
db_index=True,
)
first_event = models.DateTimeField(null=True)
external_id = models.CharField(max_length=256, null=True)
class flags(TypedClassBitField):
has_releases: bool
has_issue_alerts_targeting: bool
has_transactions: bool
has_sessions: bool
has_profiles: bool
has_replays: bool
has_feedbacks: bool
has_new_feedbacks: bool
has_minified_stack_trace: bool
has_cron_monitors: bool
has_cron_checkins: bool
# ...
设计要点:
FlexibleForeignKey:Sentry 的自定义外键字段,支持可选的db_constraint=False(避免大表间的 FK 约束带来的锁竞争)。flags:与 Organization 一样使用位字段来存储多个特性开关。has_releases、has_transactions、has_profiles等标志记录了该 Project 是否发送过对应类型的数据,这些标志被用作查询优化——例如跳过对未发送 Transactions 的 Project 的 Performance 查询。first_event:与date_added区分开来,记录 Project 真正接收到第一条事件的时间。
13.2.3 Release 模型
Release 模型记录了代码发布版本信息,定义在 src/sentry/models/release.py:208:
class Release(Model):
"""
A release is generally created when a new version is pushed into a
production state.
"""
organization = FlexibleForeignKey("sentry.Organization")
projects = models.ManyToManyField(
"sentry.Project", related_name="releases", through=ReleaseProject
)
status = BoundedPositiveIntegerField(
default=ReleaseStatus.OPEN,
choices=((ReleaseStatus.OPEN, _("Open")), (ReleaseStatus.ARCHIVED, _("Archived"))),
)
version = models.CharField(max_length=250)
ref = models.CharField(max_length=250, null=True, blank=True)
url = models.URLField(null=True, blank=True)
date_added = models.DateTimeField(default=timezone.now, db_index=True)
date_released = models.DateTimeField(null=True, blank=True)
data = LegacyTextJSONField(default=dict)
owner_id = HybridCloudForeignKey("sentry.User", on_delete="SET_NULL", null=True, blank=True)
# 物化统计
commit_count = BoundedPositiveIntegerField(null=True, default=0)
last_commit_id = BoundedBigIntegerField(null=True)
authors = ArrayField(models.TextField(), default=list, null=True)
total_deploys = BoundedPositiveIntegerField(null=True, default=0)
# 语义化版本反范式列
package = models.TextField(null=True)
major = models.BigIntegerField(null=True)
minor = models.BigIntegerField(null=True)
patch = models.BigIntegerField(null=True)
revision = models.BigIntegerField(null=True)
prerelease = models.TextField(null=True)
build_code = models.TextField(null=True)
build_number = models.BigIntegerField(null=True)
设计要点:
- 语义化版本(Semver)反范式:
version字段存储原始版本字符串(如myapp@2.14.3-beta),同时将解析后的major、minor、patch、prerelease等字段反范式化存储,以便 SQL 层高效按版本排序和过滤。 - 物化统计:
commit_count、authors、total_deploys等字段在写入时计算并存储,避免每次查询时 JOIN 计算。 HybridCloudForeignKey:owner_id引用的是通过 RPC 访问的远程 User 服务,不建立数据库级别的外键约束。- 唯一约束:
unique_together = (("organization", "version"),),确保同一组织内版本号唯一。 - PostgreSQL 函数索引:
sentry_release_semver_by_package_idx使用F("major").desc()等表达式创建函数索引,支持高效语义化版本排序。
13.2.4 Team 与 Environment
Team 模型(位于 src/sentry/models/team.py)表示组织内的团队,通过 ProjectTeam 中间表与 Project 建立多对多关系。Team 与 Organization 是一对多关系。
Environment 模型(位于 src/sentry/models/environment.py)表示部署环境(如 production、staging),与 Organization 关联,并与 Release 通过 ReleaseEnvironment 中间表建立关系。
这些模型使用标准的 Django 外键关系,相对轻量,但在权限校验和数据隔离中起到关键的过滤作用——几乎所有查询都需要通过 project__organization=org 的链路进行组织级别的数据隔离。
13.3 Group(Issue)模型详解
Group 是 Sentry 中最核心的模型之一,代表聚合后的 Issue。每一条 Group 记录汇总了一类相同根因的错误事件。定义在 src/sentry/models/group.py:848,底层数据库表名为 sentry_groupedmessage。
13.3.1 核心字段
@cell_silo_model
class Group(Model):
"""
Aggregated message which summarizes a set of Events.
"""
project = FlexibleForeignKey("sentry.Project")
logger = models.CharField(max_length=64, blank=True, default="", db_index=True)
level = BoundedPositiveIntegerField(
choices=[(key, str(val)) for key, val in sorted(LOG_LEVELS.items())],
default=logging.ERROR, blank=True, db_index=True,
)
message = models.TextField()
culprit = models.CharField(max_length=200, blank=True, null=True, db_column="view")
num_comments = BoundedPositiveIntegerField(default=0, null=True)
platform = models.CharField(max_length=64, null=True)
status = BoundedPositiveIntegerField(default=GroupStatus.UNRESOLVED, db_index=True,
choices=(
(GroupStatus.UNRESOLVED, _("Unresolved")),
(GroupStatus.RESOLVED, _("Resolved")),
(GroupStatus.IGNORED, _("Ignored")),
))
substatus = BoundedIntegerField(null=True,
choices=(
(GroupSubStatus.UNTIL_ESCALATING, _("Until escalating")),
(GroupSubStatus.ONGOING, _("Ongoing")),
(GroupSubStatus.ESCALATING, _("Escalating")),
(GroupSubStatus.UNTIL_CONDITION_MET, _("Until condition met")),
(GroupSubStatus.FOREVER, _("Forever")),
))
times_seen = BoundedPositiveIntegerField(default=1, db_index=True)
last_seen = models.DateTimeField(default=timezone.now, db_index=True)
first_seen = models.DateTimeField(default=timezone.now, db_index=True)
first_release = FlexibleForeignKey("sentry.Release", null=True, on_delete=models.PROTECT)
resolved_at = models.DateTimeField(null=True, db_index=True)
active_at = models.DateTimeField(null=True, db_index=True)
time_spent_total = BoundedIntegerField(default=0)
time_spent_count = BoundedIntegerField(default=0)
data = LegacyTextJSONField(null=True)
short_id = BoundedBigIntegerField(null=True)
type = BoundedPositiveIntegerField(default=1, db_default=1, db_index=True)
priority = models.PositiveIntegerField(null=True)
priority_locked_at = models.DateTimeField(null=True)
seer_fixability_score = models.FloatField(null=True)
seer_autofix_last_triggered = models.DateTimeField(null=True)
字段分析:
| 字段 | 类型 | 说明 |
|---|---|---|
project |
FK | 关联到 Project,数据隔离的基础 |
logger |
Char(64) | 记录日志器名称,如 sentry.errors |
level |
Int | 严重级别(error/fatal/warning等),使用 LOG_LEVELS 映射 |
message |
Text | Issue 的标题/摘要,save() 中自动截断为 255 字符 |
culprit |
Char(200) | 错误来源(db 列名为 view,历史遗留) |
platform |
Char(64) | 平台标识(python、javascript、java 等) |
status |
Int | 状态:0=未解决,1=已解决,2=已忽略 |
substatus |
Int | 子状态:ongoing、escalating、regressed 等 |
times_seen |
Int | 已统计的发生次数(计数器) |
first_seen / last_seen |
DateTime | 首次/最近发生时间 |
first_release |
FK | 首次出现的 Release 版本 |
short_id |
BigInt | Project 内的短 ID(如 PROJECT-32) |
type |
Int | Issue 类型 ID,映射到 GroupCategory(ERROR/PERFORMANCE) |
priority |
Int | 优先级分数,用于 Issue 排序 |
active_at |
DateTime | 活跃时间,默认为 first_seen |
data |
JSON | 存储额外的元数据(如完整异常的 metadata) |
BoundedPositiveIntegerField / BoundedBigIntegerField 是 Sentry 的自定义字段类型,基于 PostgreSQL 的特定整数范围(如 int4),防止了 Django 默认 AutoField 溢出问题和自增 ID 猜测。
13.3.2 数据库表结构与索引
class Meta:
app_label = "sentry"
db_table = "sentry_groupedmessage"
verbose_name_plural = _("grouped messages")
permissions = (("can_view", "Can view"),)
indexes = [
models.Index(fields=("project", "first_release")),
models.Index(fields=("project", "id")),
models.Index(fields=("project", "status", "last_seen", "id")),
models.Index(fields=("project", "status", "type", "last_seen", "id")),
models.Index(fields=("project", "status", "substatus", "last_seen", "id")),
models.Index(fields=("project", "status", "substatus", "type", "last_seen", "id")),
models.Index(fields=("project", "status", "substatus", "id")),
models.Index(fields=("status", "substatus", "id")),
models.Index(fields=("status", "substatus", "first_seen")),
models.Index(fields=("project", "status", "priority", "last_seen", "id")),
]
unique_together = (("project", "short_id"),)
索引策略分析:
索引以 project 为前缀的设计模式((project, status, last_seen, id))契合了 Sentry 的核心查询模式——几乎所有 Issue 查询都限定在 Project 范围内。这种”前缀索引”使得 PostgreSQL 可以利用 Index-Only Scan 高效获取 Issue 列表。
多个索引覆盖不同的 status + substatus + type 组合,是因为 Sentry 的 Issue Stream 页面需要根据这些维度组合过滤和排序。例如:
(project, status, last_seen, id)用于按最近活动时间列出未解决的 Issue(project, status, substatus, type, last_seen, id)用于 Performance Issue 的特定查询(project, status, priority, last_seen, id)用于按优先级排序的 Issue 流
(project, short_id) 唯一约束保证了短 ID(如 PROJECT-1)在 Project 级别唯一。
13.3.3 状态机
Group 的状态由 status 和 substatus 两个字段联合定义:
UNRESOLVED(0)
/ | \
/ | \
ONGOING ESCALATING REGRESSED NEW
|
v
RESOLVED(1) ----------> UNRESOLVED (reopen)
<---------
|
v
IGNORED(2)
/ | \
UNTIL_ESCALATING FOREVER UNTIL_CONDITION_MET
状态转换的核心方法 GroupManager.update_group_status() (group.py:695):
def update_group_status(self, groups, status, substatus, activity_type, ...):
selected_groups = Group.objects.filter(id__in=[g.id for g in groups]).exclude(
status=status, substatus=substatus
)
for group in selected_groups:
group.status = status
group.substatus = substatus
if status == GroupStatus.RESOLVED:
group.resolved_at = timezone.now()
Group.objects.bulk_update(modified_groups_list,
["status", "substatus", "priority", "resolved_at"])
for group in modified_groups_list:
Activity.objects.create_group_activity(group, activity_type, ...)
record_group_history_from_activity_type(group, activity_type.value)
关键行为:
- 批量更新(
bulk_update):一次 SQL 更新多个 Group 的状态,避免 N+1 问题。 - Activity 记录:每次状态变更创建一条 Activity 记录,用于 Issue 详情页的时间线展示。
- GroupHistory 记录:通过
record_group_history_from_activity_type()写入状态变更历史。 - Open Period 管理:通过
update_group_open_period()记录 Issue 打开/关闭的时间窗口。
GroupStatus 类定义了所有可能的宏观状态:
class GroupStatus:
UNRESOLVED = 0
RESOLVED = 1
IGNORED = 2
PENDING_DELETION = 3
DELETION_IN_PROGRESS = 4
PENDING_MERGE = 5
REPROCESSING = 6
# MUTED 是 IGNORED 的历史别名
13.3.4 GroupManager 与关联查询
GroupManager (group.py:513) 提供了丰富的查询方法:
Short ID 查询 (by_qualified_short_id_bulk):
def by_qualified_short_id_bulk(self, organization_id, short_ids_raw, *, project_ids):
# 解析 PROJECT-32 格式
# 首先按 project_slug 和 short_id 精确查询
# 不命中时回退到 case-insensitive 的 slug 查询(处理历史遗留大小写问题)
base_group_queryset = self.exclude(
status__in=[GroupStatus.PENDING_DELETION, ...]
).filter(project__organization=organization_id)
# ...
Short ID 使用了 base32 编码(通过 base32_encode/base32_decode),将数字 ID 编码为更短的字母数字字符串,类似 PROJECT-A3X。
通过 Event ID 反查 Group (from_event_id):
def from_event_id(self, project, event_id):
event = eventstore.backend.get_event_by_id(project.id, event_id)
if event:
group_id = event.group_id
if group_id is None:
raise Group.DoesNotExist()
return self.get(id=group_id)
Group 重定向 (get_group_with_redirect):当 Group 被合并时,旧 Group ID 会记录到 GroupRedirect 表。查询时自动检查重定向表,透明地返回合并后的新 Group:
def get_group_with_redirect(id_or_qualified_short_id, queryset=None, organization=None):
try:
return getter(**params), False
except Group.DoesNotExist as error:
# 查询 GroupRedirect 表
params = {
"id__in": GroupRedirect.objects.filter(
previous_group_id=id_or_qualified_short_id,
).values_list("group_id", flat=True)[:1]
}
return updated_qs.get(**params), True # True 表示发生了重定向
批量获取最新事件 (bulk_get_latest_event_ids):该函数直接从 Snuba 查询每个 Group 的最新事件 ID,通过 argMax 聚合函数获取每个 group_id 下最新的 event_id:
def bulk_get_latest_event_ids(groups):
# 按 (organization_id, dataset) 分区
for (organization_id, dataset), partition in partitions.items():
request = Request(
dataset=dataset.value,
query=Query(
match=Entity(dataset.value),
select=[
Column("group_id"),
Function("argMax", [
Column("event_id"),
Function("tuple", [Column("timestamp"), Column("event_id")]),
], "event_id"),
],
groupby=[Column("group_id")],
# ...
),
tenant_ids={"organization_id": organization_id},
)
13.4 Nodestore
Nodestore 是 Sentry 的事件体(Event Body)键值存储,专门存储单个事件的完整 JSON 负载,与 PostgreSQL 和 ClickHouse 形成互补。
13.4.1 设计理念
核心代码位于 src/sentry/services/nodestore/base.py 的 NodeStorage 类。其设计原则是:
- KV 抽象:以字符串 ID 为键,以 JSON 可序列化对象为值。
- 与 PostgreSQL 分离:事件体(可达数十 KB)不存入 PostgreSQL 主表,避免拖慢 OLTP 性能。
- 多后端可插拔:支持 Django PG 后端(小规模部署)和 Bigtable 后端(大规模部署)。
- 压缩存储:数据在写入前压缩(Django 后端用 zlib,Bigtable 后端可选 zlib/zstd)。
- Subkey 机制:允许在同一个 key 下存储多个相关的 JSON 负载,共同压缩以提升压缩比。
13.4.2 后端实现
Django 后端 (src/sentry/services/nodestore/django/backend.py):
class DjangoNodeStorage(NodeStorage):
def _get_bytes(self, id: str) -> bytes | None:
try:
data = Node.objects.get(id=id).data
return decompress(data)
except Node.DoesNotExist:
return None
def _set_bytes(self, id, data, ttl=None):
Node.objects.update_or_create(
id=id, defaults={"data": compress(data), "timestamp": timezone.now()}
)
def cleanup(self, cutoff_timestamp):
from sentry.db.deletion import BulkDeleteQuery
total_seconds = (timezone.now() - cutoff_timestamp).total_seconds()
days = math.floor(total_seconds / 86400)
BulkDeleteQuery(model=Node, dtfield="timestamp", days=days).execute()
底层 Node 模型(src/sentry/services/nodestore/django/models.py)结构简单:
class Node(BaseModel):
id = models.CharField(max_length=40, primary_key=True)
data = models.TextField() # 以 pickle 或压缩格式存储
timestamp = models.DateTimeField(default=timezone.now, db_index=True)
# db_table = "nodestore_node"
注意:Django 后端的 data 字段以 pickle 序列化存储(非 JSON),这在早期的 Sentry 设计中有历史原因,但存在安全风险——pickle 反序列化可能执行任意代码。在生产环境中,Sentry 推荐使用 Bigtable 后端。
Bigtable 后端 (src/sentry/services/nodestore/bigtable/backend.py):
class BigtableNodeStorage(NodeStorage):
store_class = BigtableKVStorage
def __init__(self, project=None, instance="sentry", table="nodestore",
automatic_expiry=False, default_ttl=None,
compression=False, **client_options):
self.store = self.store_class(
project=project, instance=instance, table_name=table,
default_ttl=default_ttl, compression=_compression,
)
self.skip_deletes = automatic_expiry and "_SENTRY_CLEANUP" in os.environ
def _get_bytes(self, id: str) -> bytes | None:
return self.store.get(id)
def _set_bytes(self, id, data, ttl=None):
self.store.set(id, data, ttl)
def delete(self, id: str):
if self.skip_deletes:
return
self.store.delete(id)
Bigtable 后端的优势:
- 自动过期(
automatic_expiry):利用 Bigtable 的 TTL 垃圾回收机制,无需显式删除。 - 无限扩展:随节点增加线性扩展。
skip_deletes:当 Bigtable 自动 GC 时(环境变量_SENTRY_CLEANUP),跳过显式删除操作以减少开销。
13.4.3 Redis 缓存层
Nodestore 内置了基于 Redis 的读取缓存层:
class NodeStorage(local, Service):
@cached_property
def cache(self) -> BaseCache | None:
try:
return caches["nodedata"]
except InvalidCacheBackendError:
return None
def get(self, id, subkey=None):
if subkey is None:
item_from_cache = self._get_cache_item(id)
if item_from_cache:
return item_from_cache
bytes_data = self._get_bytes(id)
rv = self._decode(bytes_data, subkey=subkey)
if subkey is None:
self._set_cache_item(id, rv) # 写入缓存
return rv
def _get_cache_item(self, item_id):
if self.cache:
return self.cache.get(item_id)
return None
def _set_cache_item(self, item_id, data):
if self.cache and data:
self.cache.set(item_id, data, timeout=options.get("nodestore.cache-ttl"))
缓存流程:
- 读取时优先查询 Redis(
nodedata缓存后端)。 - 缓存命中直接返回,打点
nodestore.get{cache:hit}。 - 缓存未命中则查询底层存储,解码后回写缓存。
- 写入/删除操作同时更新缓存。
- 缓存 TTL 由
nodestore.cache-ttl配置项控制。
这个缓存层对于 Issue 详情页的性能至关重要——用户在列表中反复切换 Issue 时,事件体数据通过 Redis 缓存避免了重复的 PostgreSQL/Bigtable 查询。
13.4.4 Subkey 机制
Nodestore 支持在同一个 key 下存储多个子键值,这一机制主要用于事件重处理(Reprocessing)场景:
def set_subkeys(self, item_id, data, ttl=None):
"""
>>> nodestore.set_subkeys('key1', {
... None: {'foo': 'bar'},
... "reprocessing": {'foo': 'bam'},
... })
"""
cache_item = data.get(None)
bytes_data = self._encode(data)
self.set_bytes(item_id, bytes_data, ttl=ttl)
self._set_cache_item(item_id, cache_item)
None 键存储主事件体;"reprocessing" 等具名键存储处理过程中的中间快照。编码时所有子键值合并成一个字节流再压缩:
def _encode(self, data):
# data = {None: {...}, "reprocessing": {...}}
lines = [json_dumps(data.pop(None)).encode("utf8")]
for key, value in data.items():
lines.append(key.encode("ascii"))
lines.append(json_dumps(value).encode("utf8"))
return b"\n".join(lines)
关键洞察:将多个子键值合并后共同压缩,比单独压缩每个子键值能获得更好的压缩比,因为相同的字符串模式(如 JSON 结构)可以跨子键共享压缩字典。
13.5 Snuba / ClickHouse 查询层
Snuba 是 Sentry 的时序数据查询引擎,底层使用 ClickHouse 列式存储。所有聚合查询、趋势分析、Discover 搜索都通过 Snuba 完成。Snuba 的 Python 客户端代码位于 src/sentry/snuba/。
13.5.1 Dataset 枚举
src/sentry/snuba/dataset.py 定义了所有可查询的数据集:
class Dataset(Enum):
Events = "events" # 错误事件
Transactions = "transactions" # 性能事务
Discover = "discover" # Events + Transactions 联合查询
Outcomes = "outcomes" # 用量统计物化视图
OutcomesRaw = "outcomes_raw" # 原始用量统计
Sessions = "sessions" # 会话数据(已弃用)
Metrics = "metrics" # 发布健康度指标
PerformanceMetrics = "generic_metrics" # 通用性能指标
Replays = "replays" # 会话回放
Profiles = "profiles" # 性能剖析
IssuePlatform = "search_issues" # Issue Platform 搜索
Functions = "functions" # 函数级剖析
SpansIndexed = "spans" # Span 索引数据
EventsAnalyticsPlatform = "events_analytics_platform" # EAP
对应的 ClickHouse Entity 定义在 EntityKey 枚举中:
class EntityKey(Enum):
Events = "events"
Transactions = "transactions"
Spans = "spans"
EAPItemsSpan = "eap_items_span"
EAPItems = "eap_items"
IssuePlatform = "search_issues"
Functions = "functions"
MetricsSets = "metrics_sets"
MetricsCounters = "metrics_counters"
GenericMetricsDistributions = "generic_metrics_distributions"
# ...
13.5.2 Events 数据集
Events 数据集存储所有错误事件的索引数据。关键列映射定义在 src/sentry/snuba/events.py 的 Columns 枚举中:
class Columns(Enum):
EVENT_ID = Column(
group_name="events.event_id",
event_name="event_id",
transaction_name="event_id",
discover_name="event_id",
issue_platform_name="event_id",
alias="id",
)
GROUP_ID = Column(
group_name="events.group_id",
event_name="group_id",
transaction_name=None,
discover_name="group_id",
issue_platform_name="group_id",
alias="issue.id",
)
PROJECT_ID = Column(
group_name="events.project_id",
event_name="project_id",
transaction_name="project_id",
discover_name="project_id",
issue_platform_name="project_id",
alias="project.id",
)
Column 结构中的 group_name 字段表示该列在 ClickHouse 中的实际物理列名;*_name 字段表示在各数据集中的逻辑列名;None 表示该列在对应数据集中不可用。
查询示例——从 Snuba 获取推荐事件(get_recommended_event):
def get_recommended_event(group, conditions=None, start=None, end=None, verify_replay_exists=False):
dataset = Dataset.Events if group.issue_category == GroupCategory.ERROR else Dataset.IssuePlatform
events = eventstore.backend.get_events_snql(
organization_id=group.project.organization_id,
group_id=group.id,
start=resolved_start,
end=resolved_end,
conditions=all_conditions,
limit=limit,
orderby=EventOrdering.RECOMMENDED.value,
referrer="Group.get_helpful",
dataset=dataset,
tenant_ids={"organization_id": group.project.organization_id},
inner_limit=1000,
)
EventOrdering.RECOMMENDED 定义了推荐的排序字段,优先选取有 Session Replay、有 Sampling Trace、无 Processing Error 的事件:
class EventOrdering(Enum):
RECOMMENDED = [
"-replay.id",
"-trace.sampled",
"num_processing_errors",
"-profile.id",
"-timestamp",
"-id",
]
13.5.3 Transactions 数据集
Transactions 数据集存储性能监控的事务 Span 数据。核心查询入口在 src/sentry/snuba/transactions.py 的 query() 函数,它将请求委托给 discover.query():
def query(
selected_columns, query, snuba_params, ...,
dataset: Dataset = Dataset.Discover,
fallback_to_transactions: bool = False,
referrer: str,
) -> EventsResponse:
return discover.query(
selected_columns, query, snuba_params=snuba_params,
dataset=dataset, referrer=referrer, ...
)
Transactions 数据集的特色功能包括:
- Span 级别聚合:支持
count()、p50()、p95()、p99()等百分位函数对事务耗时进行统计。 fallback_to_transactions:当 Discover 数据集中不存在对应指标时,回退到纯 Transactions 数据集查询。- 条件支持:接受
snuba_sdk.Condition对象实现灵活的 WHERE 条件。
13.5.4 Discover 数据集
Discover 是 Sentry 的聚合搜索数据集,它是 Events 和 Transactions 的联合视图。src/sentry/snuba/discover.py 提供了完整的查询功能:
def query(
selected_columns, query, snuba_params,
equations=None, orderby=None, offset=None, limit=50,
auto_fields=False, auto_aggregations=False,
use_aggregate_conditions=False, conditions=None,
functions_acl=None, transform_alias_to_input_format=False,
sample=None, has_metrics=False, skip_tag_resolution=False,
dataset=Dataset.Discover,
referrer=None,
) -> EventsResponse:
关键能力:
- 跨数据集联合:字段可同时来自 Events 和 Transactions。
- 自定义方程式(Equations):允许用户定义如
failure_rate = count_if(transaction.status, not_equals, ok) / count()的计算字段。 - 聚合条件:支持
HAVING风格的聚合后过滤。 - 分面搜索(Facets):
get_facets()函数提供字段值的分布统计。 - 直方图(Histogram):
histogram_query()支持按数值字段分桶统计。
13.5.5 其他数据集
Profiles 数据集 (snuba/profiles.py):存储应用性能剖析数据(调用栈采样),支持函数级热点分析。
SpansIndexed 数据集:存储索引化的 Span 数据,支持 Span 级别的搜索和过滤。与 Transactions 不同,它是一个完整的 Span 存储,允许按 trace_id、span_id、span.op 等维度进行精确查询。
Replays 数据集 (snuba/replays.py):存储 Session Replay 的元数据索引,包括 replay_id、project_id、timestamp、environment 等字段,用于 Replay 列表的筛选和排序。
Functions 数据集 (snuba/functions.py):基于 Profiles 构建的函数级数据集,提供更细粒度的函数调用统计。
EventsAnalyticsPlatform (EAP):较新的分析平台数据集,通过 snuba-items Topic 摄入 TraceItem Protobuf 格式的数据。
13.5.6 SnubaQuery 订阅模型
Sentry 的告警规则(Alert Rules)通过 SnubaQuery 模型将 Snuba 查询订阅化。核心模型定义在 src/sentry/snuba/models.py:
class SnubaQuery(Model):
class Type(Enum):
ERROR = 0
PERFORMANCE = 1
CRASH_RATE = 2
environment = FlexibleForeignKey("sentry.Environment", null=True, db_constraint=False)
type = models.SmallIntegerField()
dataset = models.TextField()
query = models.TextField()
group_by = ArrayField(models.CharField(max_length=200), null=True, size=100)
aggregate = models.TextField()
time_window = models.IntegerField() # 聚合窗口(秒)
resolution = models.IntegerField() # 评估间隔(秒)
extrapolation_mode = models.IntegerField()
query_snapshot = models.JSONField(null=True)
date_added = models.DateTimeField(default=timezone.now)
class QuerySubscription(Model):
class Status(Enum):
ACTIVE = 0
CREATING = 1
UPDATING = 2
DELETING = 3
DISABLED = 4
project = FlexibleForeignKey("sentry.Project", db_constraint=False)
snuba_query = FlexibleForeignKey("sentry.SnubaQuery", related_name="subscriptions")
type = models.TextField()
status = models.SmallIntegerField(default=Status.ACTIVE.value, db_index=True)
subscription_id = models.TextField(unique=True, null=True)
date_added = models.DateTimeField(default=timezone.now)
query_extra = models.TextField(null=True)
订阅生命周期:CREATING -> ACTIVE -> UPDATING/DELETING -> 删除。
告警规则的执行依赖 Kafka 消费 Snuba 的订阅结果 Topic(如 events-subscription-results、transactions-subscription-results),由 QuerySubscriptionDataSourceHandler 管理生命周期。
13.6 Redis / Memcached 缓存策略
Sentry 的缓存体系以 Redis 为主力,在部署架构中承担多个角色的缓存任务:
1. Nodestore 事件体缓存
如 13.4.3 节所述,nodedata 缓存后端(Redis)缓存事件 JSON 体,大幅减少 Bigtable/PostgreSQL 的读取压力。TTL 由 nodestore.cache-ttl 配置。
2. Django 默认缓存(default)
Sentry 的 Django 配置默认使用 Redis 作为缓存后端,用于缓存 Django ORM 对象(如 Group、Project 配置)。部分 Manager 声明了 cache_fields 启用自动缓存:
class GroupManager(BaseManager["Group"]):
# ...
# Group 模型声明
objects: ClassVar[GroupManager] = GroupManager(cache_fields=("id",))
class QuerySubscription(Model):
objects: ClassVar[BaseManager[Self]] = BaseManager(
cache_fields=("pk", "subscription_id"),
cache_ttl=int(timedelta(hours=1).total_seconds())
)
cache_fields 指定了哪些字段可以作为缓存键。例如 Group 可以通过 id 从缓存获取,QuerySubscription 可以通过 pk 或 subscription_id 获取,TTL 为 1 小时。
3. 计数器缓存
times_seen 等计数器字段并非实时更新到 PostgreSQL——Sentry 使用 Redis 作为写入缓冲,定期将增量同步到数据库,避免高并发下的写锁竞争。
4. 选项/配置缓存(options)
Sentry 的 sentry.options 系统(类似 Django settings 的动态版本)使用 Redis 缓存配置项,支持运行时更新而无需重启。
5. 项目/组织配置缓存
通过 sentry.app.env 和 sentry.locks 模块,项目配置(如 DSN、采样率、过滤规则)被缓存在 Redis 中,由 Relay 摄取层快速读取,避免每次事件处理都查询 PostgreSQL。
6. Killswitch 缓存
killswitch_matches_context() 函数在 Kafka 消息路由中用于动态开关控制,其开关状态通过 Redis 缓存实现近实时的配置下发。
7. 查询结果缓存
某些昂贵的 Snuba 聚合查询结果可以通过 Django 的 cache API 缓存到 Redis,通过 Referrer 标识区分不同的缓存键。
13.7 Kafka 消息队列
Kafka 在 Sentry 中承担异步解耦和数据管道的角色。所有事件数据从摄取层到 ClickHouse 的传输、Post-Processing 任务的分发、告警订阅的评估结果,都通过 Kafka Topic 流转。
13.7.1 Topic 列表
完整的 Kafka Topic 枚举定义在 src/sentry/conf/types/kafka_definition.py 的 Topic 枚举中,以下为分类汇总:
| 类别 | Topic | 用途 |
|---|---|---|
| 事件流 | events |
错误事件主通道 |
transactions |
事务事件主通道 | |
generic-events |
Issue Platform 事件通道 | |
snuba-items |
EAP TraceItem 通道 | |
| Commit Log | snuba-commit-log |
Snuba events 写入确认 |
snuba-transactions-commit-log |
Snuba transactions 写入确认 | |
snuba-generic-events-commit-log |
通用事件写入确认 | |
| 摄入管道 | ingest-events |
事件摄入预处理 |
ingest-transactions |
事务摄入预处理 | |
ingest-spans |
Span 摄入 | |
ingest-metrics |
指标摄入 | |
ingest-performance-metrics |
性能指标摄入 | |
ingest-replay-events |
回放事件摄入 | |
ingest-occurrences |
Issue Occurrence 摄入 | |
ingest-attachments |
附件摄入(Minidump 等) | |
ingest-monitors |
Cron Monitor 摄入 | |
| 死信队列 | ingest-events-dlq |
事件摄入 DLQ |
ingest-transactions-dlq |
事务摄入 DLQ | |
ingest-generic-metrics-dlq |
通用指标 DLQ | |
outcomes-dlq |
用量统计 DLQ | |
| 订阅结果 | events-subscription-results |
错误告警订阅结果 |
transactions-subscription-results |
事务告警订阅结果 | |
generic-metrics-subscription-results |
通用指标告警结果 | |
metrics-subscription-results |
发布健康度告警结果 | |
subscription-results-eap-items |
EAP Items 告警结果 | |
| 用量统计 | outcomes |
用量物化视图 |
outcomes-billing |
计费用量数据 | |
| 会话回放 | ingest-replay-recordings |
回放录制数据 |
| 性能 | profiles |
性能剖析数据 |
processed-profiles |
已处理的剖析数据 | |
snuba-profile-chunks |
剖析块数据 | |
| 监控 | monitors-clock-tick |
Monitor 时钟滴答 |
monitors-incident-occurrences |
Monitor 事件 | |
uptime-results |
Uptime 检查结果 | |
| 基础设施 | group-attributes |
Group 属性变更 |
taskworker |
异步任务工作 | |
shared-resources-usage |
共享资源用量 | |
buffered-segments |
缓冲分段数据 |
13.7.2 EventStream 消息格式
Sentry 的 EventStream 系统有两个版本的消息协议(见 src/sentry/eventstream/snuba.py:37-82):
Version 1 格式:
(1, TYPE, [...REST...])
# Insert: (1, 'insert', {event json}, {post-process state})
# 仅 insert 操作被处理
Version 2 格式(当前版本):
(2, TYPE, [...REST...])
# Insert: (2, 'insert', {event json}, {post-process state})
# Delete Groups: (2, '(start_delete_groups|end_delete_groups)', {...})
# Merge: (2, '(start_merge|end_merge)', {...})
# Unmerge: (2, '(start_unmerge|end_unmerge)', {...})
# Delete Tag: (2, '(start_delete_tag|end_delete_tag)', {...})
Insert 消息中的 event JSON 包含完整的规范化事件数据:
extra_data = (
{
"group_id": event.group_id,
"group_ids": [group.id for group in getattr(event, "groups", [])],
"event_id": event.event_id,
"organization_id": project.organization_id,
"project_id": event.project_id,
"message": event.search_message,
"platform": event.platform,
"datetime": json.datetime_to_str(event.datetime),
"data": event_data,
"primary_hash": primary_hash,
"retention_days": retention_days,
"occurrence_id": occurrence_data.get("id"),
"occurrence_data": occurrence_data,
},
{
"is_new": is_new,
"is_regression": is_regression,
"is_new_group_environment": is_new_group_environment,
"skip_consume": skip_consume,
"group_states": group_states,
},
)
protocol.py 中的消息处理器提取 Post-Process 任务所需参数:
def get_task_kwargs_for_insert(operation, event_data, task_state):
if task_state and task_state.get("skip_consume", False):
return None
kwargs = {
"event_id": event_data["event_id"],
"project_id": event_data["project_id"],
"group_id": event_data["group_id"],
"primary_hash": event_data["primary_hash"],
"occurrence_id": event_data.get("occurrence_id"),
}
for name in ("is_new", "is_regression", "is_new_group_environment"):
kwargs[name] = task_state[name]
return kwargs
13.7.3 消息发送流程
KafkaEventStream (src/sentry/eventstream/kafka/backend.py) 实现了完整的 Kafka 消息发送:
class KafkaEventStream(SnubaProtocolEventStream):
def __init__(self, **options):
super().__init__(**options)
self.topic = Topic.EVENTS
self.transactions_topic = Topic.TRANSACTIONS
self.issue_platform_topic = Topic.EVENTSTREAM_GENERIC
def _send(self, project_id, _type, extra_data=(), asynchronous=True,
headers=None, skip_semantic_partitioning=False,
event_type=EventStreamEventType.Error):
if event_type == EventStreamEventType.Transaction:
topic = self.get_transactions_topic(project_id)
elif event_type == EventStreamEventType.Generic:
topic = self.issue_platform_topic
else:
topic = self.topic
producer = self.get_producer(topic)
real_topic = get_topic_definition(topic)["real_topic_name"]
produce_future = producer.produce(
destination=ArroyoTopic(real_topic),
payload=KafkaPayload(
key=str(project_id).encode("utf-8") if not skip_semantic_partitioning else None,
value=json.dumps((self.EVENT_PROTOCOL_VERSION, _type) + extra_data).encode("utf-8"),
headers=[(k, v.encode("utf-8")) for k, v in headers.items()],
),
)
分区策略:
- 语义分区(Semantic Partitioning):默认以
project_id作为 Kafka Key。同一 Project 的事件进入同一分区,保证同一个 Project 内的时序一致性。 - 随机分区:对于 Transactions 和 Generic Events,以及被 Killswitch 标记的 Project,使用
skip_semantic_partitioning=True,以None作为 Key 实现随机分布,避免热点分区问题。 - 同步/异步:默认异步发送(
asynchronous=True),采样事件(sample_event标签)和 Delete/Merge 等控制消息使用同步发送以确保可靠性。
Post-Process Forwarder:KafkaEventStream.requires_post_process_forwarder() 返回 True(而直连 Snuba 的 SnubaEventStream 返回 False),表示需要一个独立的 Forwarder 进程从 Kafka 消费事件消息并触发 post_process_group 任务(如发送通知、处理规则匹配)。
消息头优化:当 eventstream:kafka-headers 选项启用时,关键的元数据(event_id、project_id、group_id、is_new 等)直接放入 Kafka 消息头,Post-Process Forwarder 无需解析完整的 JSON 消息体即可做出路由决策,降低了单核 CPU 的解析开销。
13.8 数据库迁移管理
Sentry 使用 Django 的 Migration 框架管理 PostgreSQL Schema 变更。迁移文件位于各应用目录下的 migrations/ 子目录中。
关键实践:
BulkDeleteQuery:Sentry 实现了自定义的批量删除查询类,用于大规模数据清理操作(如 Nodestore 清理、旧事件删除)。它使用分批次 + 游标的方式,避免长时间锁定表。
# Nodestore Django 后端的清理
def cleanup(self, cutoff_timestamp):
from sentry.db.deletion import BulkDeleteQuery
total_seconds = (timezone.now() - cutoff_timestamp).total_seconds()
days = math.floor(total_seconds / 86400)
BulkDeleteQuery(model=Node, dtfield="timestamp", days=days).execute()
- 零停机迁移策略:
- 在线 Schema 变更:利用 PostgreSQL 的
CREATE INDEX CONCURRENTLY创建索引而不阻塞写入。 - 分阶段迁移:复杂变更拆分为多步——先添加新列(可为 NULL),再填充数据,最后添加 NOT NULL 约束。
- 原子操作:每个 Migration 文件包含一个事务性操作,确保变更要么全部完成,要么全部回滚。
-
Snowflake ID 迁移:对于 Organization 等模型,通过
save_with_snowflake_id()函数在保存时生成 Snowflake ID,避免依赖数据库自增序列。 -
测试事务:
in_test_hide_transaction_boundary()上下文管理器在测试中隐藏事务边界,使迁移的测试更加可靠。 -
cell_silo_model:标记了cell_silo_model的模型受 Silo 架构约束,其迁移必须考虑跨 Region 的数据一致性和复制影响。
13.9 数据生命周期管理
Sentry 的数据生命周期管理涉及多个层面的 TTL 和清理策略:
1. 事件保留期(Event Retention)
每个组织可以配置事件保留天数,由 quotas.backend.get_event_retention() 获取。在查询 Snuba 时,outside_retention_with_modified_start() 函数会检查请求的时间范围是否超出保留期:
expired, start = outside_retention_with_modified_start(
start, end, Organization(organization_id)
)
if expired:
return None # 数据已过期,跳过查询
2. Nodestore 清理
- Django 后端:通过
cleanup()方法定期调用BulkDeleteQuery删除过期的 Node 记录。 - Bigtable 后端:可选启用
automatic_expiry,由 Bigtable 的列族 TTL 自动删除过期数据,同时设置skip_deletes=True跳过显式删除。
3. ClickHouse TTL
ClickHouse 的 MergeTree 表引擎原生支持 TTL(Time To Live),可以通过 TTL timestamp + INTERVAL 90 DAY DELETE 在表定义层面自动删除过期数据。Snuba 的 ClickHouse Schema 中为各数据集配置了不同的 TTL。
4. Group 删除流程
Group 的删除是一个多阶段流程:
- 状态先变为
PENDING_DELETION(3)。 - 进入
DELETION_IN_PROGRESS(4),此时禁止新事件关联到此 Group。 - 通过 EventStream 发送
tombstone_events和exclude_groups消息到 Snuba。 - Snuba 在 ClickHouse 侧标记事件为已删除,并从查询结果中排除该 Group。
- 最后在 PostgreSQL 中物理删除 Group 记录。
5. 归档策略
对于已解决的 Issue,Sentry 不会立即删除其数据,而是保留完整的查询能力直到保留期结束。IGNORED 状态的 Issue 有特殊的子状态:
UNTIL_ESCALATING:保持忽略直到事件数量达到阈值重新激活UNTIL_CONDITION_MET:直到满足特定条件重新激活FOREVER:永久忽略
6. 数据采样
Sentry 的 Discover 查询支持 sample 参数,当数据量过大时可以按比例采样以减少查询负载。这在大型组织的 Discover 自定义查询中用于提高响应速度。
13.10 存储选择指南
以下表格汇总了各类数据应存储在哪里以及选择理由:
| 数据类型 | 存储位置 | 理由 |
|---|---|---|
| 组织/项目/团队/用户等元数据 | PostgreSQL | 需要事务一致性、关系查询(JOIN)、权限校验 |
| Issue 元数据(状态、计数、时间) | PostgreSQL(sentry_groupedmessage) |
需要行级更新、分页查询、状态流转 |
| 事件完整 JSON 体 | Nodestore(PostgreSQL/Bigtable) | 数十 KB 的大字段不适合 PG 主表;KV 模型简化存取 |
| 事件索引(project_id, group_id, timestamp, tags…) | ClickHouse (Snuba) | 列存适合聚合查询、高压缩比、时间范围扫描 |
| 性能事务数据 | ClickHouse (Snuba Transactions) | 需要百分位聚合(p50/p95/p99)、高写入吞吐 |
| 性能剖析数据(Profiles) | ClickHouse (Snuba Profiles) | 函数级聚合、调用栈分析 |
| Span 索引数据 | ClickHouse (Snuba SpansIndexed) | Span 级别搜索、按 trace_id 追踪 |
| 会话回放元数据 | ClickHouse (Snuba Replays) | 回放列表筛选、时间范围查询 |
| 告警规则与订阅 | PostgreSQL(sentry_snubaquery / sentry_querysubscription) |
需事务性的创建/更新/删除 |
| 告警评估结果 | Kafka -> PostgreSQL (Incident) | 通过 Kafka *-subscription-results 传输,触发 Incident 创建 |
| 用量统计数据 | ClickHouse (Snuba Outcomes) | 物化视图聚合、计费报表 |
| 事件流缓冲 | Kafka | 削峰填谷、异步解耦、消费者可独立扩展 |
| 事件体缓存 | Redis (nodedata) |
亚毫秒读取、减轻 Nodestore 负载 |
| 项目配置缓存 | Redis | Relay 需要低延迟读取 DSN/过滤规则/采样率 |
| 查询结果缓存 | Redis | 避免重复的昂贵 Snuba 聚合查询 |
| 计数器增量缓冲 | Redis | 高并发写入缓冲,定期刷入 PG |
核心原则:
- “元数据在 PG,事件体在 Nodestore,索引在 ClickHouse” —— 这是 Sentry 存储架构的核心分界。
- PostgreSQL 是真相源(Source of Truth) —— 所有可变的业务状态数据以 PG 为准,Snuba 和 Nodestore 可从 PG 重新构建。
- ClickHouse 是查询加速器 —— 不做 OLTP 行级操作,只做 OLAP 聚合查询。数据最终一致性可接受。
- Kafka 是缓冲带 —— 解耦写入路径和消费路径。摄入层快速确认(200 OK),消息可靠持久化后异步处理。
- Redis 是速度层 —— 缓存一切热点数据,命中率直接决定页面加载延迟。
- 冗余是设计的一部分 —— 同一个事件在 Nodestore(完整 JSON)、PostgreSQL(Group 聚合)、ClickHouse(索引 + Tags)三处各有副本,各司其职。