Sentry Replays 录制消费端消息契约解析:从 recording_consumer Blueprint 到 ingest-replay-recordings 消费实现
2026/9/10 2:23:53 网站建设 项目流程

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_chunkedreplay_recording_chunkreplay_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):

字段类型说明
typestring字面量:replay_recording_not_chunked
replay_idstring回放会话的唯一 ID
key_idOptional[int]上报所用的 API Key ID(可选)
org_idint组织 ID
project_idint项目 ID
receivedint事件被接收时的 Unix 时间戳
payloadbytesJSON 编码的头部信息会被前缀到 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(录制分块)

字段类型说明
typestring字面量:replay_recording_chunk
replay_idstring回放会话的唯一 ID
project_idint项目 ID
chunk_indexint分块序号(从 0 开始)
idstring唯一标识(关联汇总消息)
payloadbytes该块的录制字节

示例:

{ "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_idorg_idreceived字段,也没有JSON 头部前缀(payload 直接是分块字节);id字段用于把多个 chunk 关联到同一条录制汇总消息上。

4.2 Replay Recording(录制汇总)

字段类型说明
typestring字面量:replay_recording
replay_idstring回放会话的唯一 ID
key_idOptional[int]上报所用的 API Key ID(可选)
org_idint组织 ID
project_idint项目 ID
receivedint事件被接收时的 Unix 时间戳
replay_recordingdict包含分块数量(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_chunkedchunkrecording(汇总)
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_asyncdelay调用它,因此任务签名、名称与命名空间不可随意变更。

任务内部两个阶段清晰分离:

阶段一:处理(Processing Task)——process_message 依次执行:

  1. parse_request_message:用RECORDINGS_CODEC做 msgpack 解码并做 schema 校验;
  2. parse_headers:按\n拆分 JSON 头与分段体,得到segment_id
  3. decompress_segment:zlib 解压(含未压缩兜底);
  4. 提取可选的replay_event(replay 事件 JSON)与replay_video(视频字节),并读取relay_snuba_publish_disabled标记决定是否由本消费端代发 replay event 到 Snuba 消费端(见 parse_recording_event);
  5. 调用 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:

  1. 写入对象存储storage_kv.set(recording.filename, recording.filedata),即把(压缩后的)录制分段写入 KV 存储;
  2. 计费与首段标记:当segment_id == 0(首个分段)时,通过_track_initial_segment_event_new/_old对项目设置has_replays标志、触发first_replay_received信号,并写入DataCategory.REPLAY的 ACCEPTED outcome 计费记录;
  3. 按需补发 replay event:若relay_snuba_publish_disabled为 True 且携带了replay_event,则通过publish_replay_event(见 kafka.py)把事件发往ingest-replay-eventstopic,交给 Snuba 消费端;
  4. 事件落库与 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)**划定了稳定的接口边界。基于文档与源码,可以总结出以下排障与二次开发要点:

  1. 编码必须匹配:所有消息必须是 msgpack 字节(msgpack.packb),否则RECORDINGS_CODEC.decode会抛ValidationError并被静默丢弃(测试中以空消息{}验证了这一行为)。
  2. 1MB 阈值决定消息形态:小于 1MB 用replay_recording_not_chunked单条投递;大于等于 1MB 必须拆分为replay_recording_chunk+replay_recording汇总,且汇总中的id与各 chunk 的id必须一致、chunks必须等于实际块数。
  3. payload 首行必须是 JSON 头replay_recording_not_chunked的 payload 以{"segment_id": N}+\n开头,解析依赖首个\n切分;缺失或格式错误会被parse_headers静默丢弃。
  4. 解压是自适应的:消费端同时接受 zlib 压缩与非压缩(JSON 数组开头)的分段体,SDK/Relay 侧可按需压缩。
  5. 字段差异要分清: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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询