Sentry Replays 录制消费端消息契约解析:从 recording_consumer Blueprint 到 ingest-replay-recordings 消费实现
【免费下载链接】sentryDeveloper-first error tracking and performance monitoring项目地址: https://gitcode.com/GitHub_Trending/sen/sentry
Sentry Replays 的录制(recording)数据从浏览器 SDK 采集后,经由 Relay 进入后端消息管道,最终由"录制消费端"(recording consumer)负责解码、解压、解析并落库。本文以仓库中 recording_consumer.md 这份 Blueprint 文档为主体,完整讲解录制消费端与其生产者之间的消息契约——包括小录制(Small Recordings)与大录制(Large Recordings)两类消息的字段定义、msgpack 编码约定、1MB 分块阈值,并结合 recording.py、tasks.py、ingest/init.py 等源码,还原从"收到字节"到"写入对象存储与计费"的完整链路,帮助你掌握自建 Relay/SDK 对接 Sentry Replays 时的消息格式规范与排查依据。
一、契约文档定位:消费端与生产者的"接口约定"
src/sentry/replays/blueprints/目录下存放的是 Sentry Replays 各条数据管线的接口约定文档,与 recording_consumer.md 并列的还有 relay.md(SDK 事件与录制分段事件契约)与 snuba_consumer.md(写入 Snuba 的事件契约)。其中 recording_consumer.md 明确说明了它的定位:
This file defines the contract between the recording consumer and its producers.
即:这份文档定义的是录制消费端与它的生产者(Relay 以及 SDK 产生的分段数据)之间的消息契约。理解这份契约,需要先厘清几个贯穿全文的关键约定:
- 编码格式:所有消息都以
msgpack编码的字节流(bytes)传输,而不是 JSON 文本。 - 1MB 分块阈值:消息体小于 1MB 时,作为单条消息整体处理,其结构定义在 "Small Recordings"(小录制)一节;其余所有消息类型定义在 "Large Recordings"(大录制)一节。
- 消息类型由
type字段区分,共有三种:replay_recording_not_chunked、replay_recording_chunk、replay_recording。 - Content-Type:请求体统一为
application/octet-stream(二进制流)。
从源码侧看,这份契约在 Kafka topic 层面有对应实现。src/sentry/conf/types/kafka_definition.py中定义了主题枚举INGEST_REPLAYS_RECORDINGS = "ingest-replay-recordings"(见 kafka_definition.py),录制消费端正是消费该 topic 的消息。而 recording.py 通过get_topic_codec(Topic.INGEST_REPLAYS_RECORDINGS)拿到对应的 msgpack Codec,用于消息解码。
二、消息传输与编码:msgpack、1MB 阈值与 Codec 校验
2.1 为什么用 msgpack
Blueprint 开篇强调所有消息均以msgpack编码。与 JSON 相比,msgpack 是二进制序列化格式,体积更小、解析更快,非常适合 Replays 这种以高吞吐字节流为主的录制数据链路。
在消费端,消息解码由 parse_request_message 完成:
@trace def parse_request_message(message: bytes) -> ReplayRecording: try: return RECORDINGS_CODEC.decode(message) except ValidationError: logger.exception("Could not decode recording message.") raise DropSilently()这里RECORDINGS_CODEC来自sentry_kafka_schemas,其 schema 类型为ingest_replay_recordings_v1。解码失败(ValidationError)时抛出DropSilently异常,消费端会静默丢弃该消息而不中断消费进度——这是录制类数据"尽力而为"处理策略的体现。
2.2 1MB 阈值的含义
文档规定"Messages smaller than 1MB in size are processed as one message",即:
- 小于 1MB 的录制:整条消息即为一个完整的录制分段,使用
replay_recording_not_chunked类型直接投递; - 大于等于 1MB 的录制:需要先被切分为多个 chunk(块),先投递多个
replay_recording_chunk消息,最后再投递一条replay_recording汇总消息声明"该录制共有多少块、唯一标识是什么",由消费端按chunk_index重组。
这种设计避免单条 Kafka 消息过大导致的分区负载不均与传输开销,是大录制场景下的标准分片策略。
2.3 测试对契约的印证
单元测试 test_recording.py 展示了契约的端到端用法:用msgpack.packb(message)对消息字典做 msgpack 编码,再调用任务入口process_replay_recording(...)处理;集成测试 integration/consumers/test_recording.py 同样使用msgpack.packb(message)构造输入。这从测试侧印证了 Blueprint 的"msgpack encoded bytes"约定。此外,集成测试还专门构造了msgpack.packb({})(空消息)验证非法消息会被静默丢弃、不触发提交(见 test_recording.py)。
三、小录制(Small Recordings):Replay Recording Not Chunked
3.1 字段契约
小录制只有一种消息类型replay_recording_not_chunked,其字段定义如下(来自 recording_consumer.md):
| 字段 | 类型 | 说明 |
|---|---|---|
| type | string | 字面量:replay_recording_not_chunked |
| replay_id | string | 回放会话的唯一 ID |
| key_id | Optional[int] | 上报所用的 API Key ID(可选) |
| org_id | int | 组织 ID |
| project_id | int | 项目 ID |
| received | int | 事件被接收时的 Unix 时间戳 |
| payload | bytes | JSON 编码的头部信息会被前缀到 payload 上 |
请求的 Content-Type 为application/octet-stream。文档给出的 Python 风格请求示例:
{ "type": "replay_recording_not_chunked", "replay_id": "515539018c9b4260a6f999572f1661ee", "key_id": 1, "org_id": 132, "project_id": 10459681, "received": 1342632621, "payload": b"{'segment_id': 0}\n\x14ftypqt\x00\x00\x00\x00qt\x00\x00x08wide\x03\xbdd\x11mdat" }3.2 payload 的内部结构:JSON 头部 + 分段体
文档特别指出"JSON encoded headers are prefixed to the payload"——即payload字段并不是纯粹的录制数据,而是"JSON 头 + 换行 + 录制分段体"。这与 relay.md 中"Replay Recording Segment Event"一节的格式一致:
{"segment_id": 0} \x00\x00\x00\x14ftypqt \x00\x00\x00\x00qt \x00\x00\x00\x08wide\x03\xbdd\x11mdat第一行{"segment_id": 0}是 JSON 编码的分段头部(声明该 payload 属于第几个 segment),换行符之后才是真正的录制数据(示例中以ftypqt/wide/mdat开头的二进制内容,是浏览器录制数据常见的 MP4/容器风格字节序列)。
消费端在 parse_headers 中精确实现了这一拆分逻辑:
@trace def parse_headers(recording: bytes, replay_id: str) -> tuple[int, bytes]: try: recording_headers_json, recording_segment = recording.split(b"\n", 1) return int(json.loads(recording_headers_json)["segment_id"]), recording_segment except Exception: logger.exception("Recording headers could not be extracted %s", replay_id) raise DropSilently()即:以第一个\n为界,前段 JSON 解析出segment_id,后段为实际的录制分段字节;任一环节失败即静默丢弃。
3.3 解压逻辑与格式兜底
分段体可能是压缩的,也可能未压缩。decompress_segment 实现了智能解压:
@trace def decompress_segment(segment: bytes) -> tuple[bytes, bytes]: try: return (segment, zlib.decompress(segment)) except zlib.error: if segment and segment[0] == ord("["): return (zlib.compress(segment), segment) else: logger.exception("Invalid recording body.") raise DropSilently()逻辑要点:
- 优先尝试
zlib.decompress,成功则返回(压缩原文, 解压结果)二元组; - 若解压失败,检查首字节是否为
[(ASCII 91,即 JSON 数组的开头字符)——如果是,说明该段本就未压缩,直接对原文做zlib.compress以便统一落库,同时返回原文作为解析用 payload; - 两者都不是,视为非法录制体,静默丢弃。
这一"首字节判 JSON"的兜底策略与 pack.py 中"过去编码以[字符开头、视为 rrweb 类型直接返回"的兼容逻辑互为印证。
四、大录制(Large Recordings):Chunk 与汇总消息
大录制场景下共两种消息类型,配合完成"分块上传 → 汇总声明"两步。
4.1 Replay Recording Chunk(录制分块)
| 字段 | 类型 | 说明 |
|---|---|---|
| type | string | 字面量:replay_recording_chunk |
| replay_id | string | 回放会话的唯一 ID |
| project_id | int | 项目 ID |
| chunk_index | int | 分块序号(从 0 开始) |
| id | string | 唯一标识(关联汇总消息) |
| payload | bytes | 该块的录制字节 |
示例:
{ "type": "replay_recording_chunk", "replay_id": "515539018c9b4260a6f999572f1661ee", "project_id": 1, "chunk_index": 10, "id": "e4a28052c54743a286be419c9d168ef5", "payload": b"\x14ftypqt\x00\x00\x00\x00qt\x00\x00x08wide\x03\xbdd\x11mdat" }注意与replay_recording_not_chunked的差异:chunk 消息不携带key_id、org_id、received字段,也没有JSON 头部前缀(payload 直接是分块字节);id字段用于把多个 chunk 关联到同一条录制汇总消息上。
4.2 Replay Recording(录制汇总)
| 字段 | 类型 | 说明 |
|---|---|---|
| type | string | 字面量:replay_recording |
| replay_id | string | 回放会话的唯一 ID |
| key_id | Optional[int] | 上报所用的 API Key ID(可选) |
| org_id | int | 组织 ID |
| project_id | int | 项目 ID |
| received | int | 事件被接收时的 Unix 时间戳 |
| replay_recording | dict | 包含分块数量(chunks)与唯一标识(id) |
示例:
{ "type": "replay_recording", "replay_id": "515539018c9b4260a6f999572f1661ee", "key_id": 1, "org_id": 132, "project_id": 10459681, "received": 1342632621, "replay_recording": { "id": "e4a28052c54743a286be419c9d168ef5", "chunks": 5 } }replay_recording子对象中的id必须与各replay_recording_chunk消息中的id一致,chunks声明总块数。消费端据此得知该录制共 5 块(示例中chunk_index取 0~4),进而按索引重组出完整录制。
4.3 三种消息的字段差异小结
| 字段 | not_chunked | chunk | recording(汇总) |
|---|---|---|---|
| type | 有 | 有 | 有 |
| replay_id | 有 | 有 | 有 |
| project_id | 有 | 有 | 有 |
| key_id | 可选 | 无 | 可选 |
| org_id | 有 | 无 | 有 |
| received | 有 | 无 | 有 |
| payload | 有(JSON 头 + 分段体) | 有(纯分块字节) | 无 |
| chunk_index | 无 | 有 | 无 |
| id | 无 | 有 | 无(在 replay_recording.id 中) |
| replay_recording | 无 | 无 | 有(含 id 与 chunks) |
五、从字节到落库:消费端完整处理链路
理解了消息契约之后,再看消费端如何消费这些消息。当前版本的 Sentry 中,ingest-replay-recordingstopic 的消费由 taskbroker 以 "raw mode" 直接投递到任务 process_replay_recording:
@instrumented_task( name=PROCESS_REPLAY_RECORDING_TASK_NAME, namespace=replays_raw_tasks, retry=Retry(times=3, delay=5), silo_mode=SiloMode.CELL, ) def process_replay_recording(message_bytes: bytes) -> None: processed_message = process_message(message_bytes) if processed_message: context: ProcessorContext = { "has_sent_replays_cache": None, "options_cache": None, } commit_message(processed_message, context)其中PROCESS_REPLAY_RECORDING_TASK_NAME = "sentry.replays.tasks.process_replay_recording"(定义于 kafka.py),且任务注释明确指出:该任务直接由 taskbroker 从ingest-replay-recordingstopic 逐条投递,应用代码不会apply_async或delay调用它,因此任务签名、名称与命名空间不可随意变更。
任务内部两个阶段清晰分离:
阶段一:处理(Processing Task)——process_message 依次执行:
parse_request_message:用RECORDINGS_CODEC做 msgpack 解码并做 schema 校验;parse_headers:按\n拆分 JSON 头与分段体,得到segment_id;decompress_segment:zlib 解压(含未压缩兜底);- 提取可选的
replay_event(replay 事件 JSON)与replay_video(视频字节),并读取relay_snuba_publish_disabled标记决定是否由本消费端代发 replay event 到 Snuba 消费端(见 parse_recording_event); - 调用 process_recording_event 生成
ProcessedEvent:解析 rrweb 事件(默认json.loads,开启replay.consumer.msgspec_recording_parser选项后改用 msgspec 快速解析,见parse_recording_data)、生成存储文件名、决定是否打包 replay video。
阶段二:提交(I/O Task)——commit_message 调用 commit_recording_message:
- 写入对象存储:
storage_kv.set(recording.filename, recording.filedata),即把(压缩后的)录制分段写入 KV 存储; - 计费与首段标记:当
segment_id == 0(首个分段)时,通过_track_initial_segment_event_new/_old对项目设置has_replays标志、触发first_replay_received信号,并写入DataCategory.REPLAY的 ACCEPTED outcome 计费记录; - 按需补发 replay event:若
relay_snuba_publish_disabled为 True 且携带了replay_event,则通过publish_replay_event(见 kafka.py)把事件发往ingest-replay-eventstopic,交给 Snuba 消费端; - 事件落库与 EAP:解析出的点击、tap、hydration error 等"高亮事件"通过
emit_replay_events分发给各事件记录器,trace items 则通过 write_trace_items 发往 EAP 的snuba-itemstopic。
存储文件名的生成规则
storage.py 中的_make_recording_filename定义了落库 Key 格式:
def _make_recording_filename(retention_days, project_id, replay_id, segment_id) -> str: return f"{retention_days or 30}/{project_id}/{replay_id}/{segment_id}"即对象存储 Key 为{retention_days}/{project_id}/{replay_id}/{segment_id},默认保留期 30 天,且 Key 前缀即为 TTL(StorageBlob注释说明支持 30/60/90 天三种保留期,需在存储桶上启用 TTL)。读取侧 API 也遵循同一命名规则,例如 api.md 中的 recording-segments 查询接口。
六、打包格式:录制与视频的二进制容器
当消息携带replay_video时,消费端会将 rrweb 录制与视频合包压缩后落库(见 pack_replay_video)。二进制容器格式定义在 pack.py:
- 首字节为 8 位类型标记:
0表示 rrweb,1表示 video; - 为兼容历史数据,类型字节最大值限定为 90(91 是 JSON 数组开头的
[的 ASCII 码,旧编码以[开头则视为未打包的 rrweb,直接返回); - video 类型布局为:
1 字节类型+4 字节长度头+video 字节+rrweb 字节,其中 4 字节长度头记录 video 部分的字节数,用于拆分(HEADER_OFFSET = 4 + 1)。
测试 test_recording.py 中使用unpack(zlib.decompress(result))[1]读取存储内容,即:先 zlib 解压,再按此容器格式拆包得到原始 rrweb JSON 数组,与写入前比对断言数据一致。
七、消息契约的工程价值与排障要点
这份 Blueprint 文档的工程价值在于:它为**生产者(Relay / SDK 分段上传器)与消费端(ingest-replay-recordings)**划定了稳定的接口边界。基于文档与源码,可以总结出以下排障与二次开发要点:
- 编码必须匹配:所有消息必须是 msgpack 字节(
msgpack.packb),否则RECORDINGS_CODEC.decode会抛ValidationError并被静默丢弃(测试中以空消息{}验证了这一行为)。 - 1MB 阈值决定消息形态:小于 1MB 用
replay_recording_not_chunked单条投递;大于等于 1MB 必须拆分为replay_recording_chunk+replay_recording汇总,且汇总中的id与各 chunk 的id必须一致、chunks必须等于实际块数。 - payload 首行必须是 JSON 头:
replay_recording_not_chunked的 payload 以{"segment_id": N}+\n开头,解析依赖首个\n切分;缺失或格式错误会被parse_headers静默丢弃。 - 解压是自适应的:消费端同时接受 zlib 压缩与非压缩(JSON 数组开头)的分段体,SDK/Relay 侧可按需压缩。
- 字段差异要分清:chunk 消息不携带
org_id/key_id/received,只有汇总消息才带完整上下文,构造消息时切勿混用。
如需进一步了解上下游,可继续阅读:上游 SDK 事件与录制分段事件契约见 relay.md;下游写入 Snuba 的事件契约见 snuba_consumer.md;消费端完整实现见 recording.py;消息处理与落库逻辑见 usecases/ingest/init.py 与 tasks.py。
【免费下载链接】sentryDeveloper-first error tracking and performance monitoring项目地址: https://gitcode.com/GitHub_Trending/sen/sentry
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考