ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Sentry Replays 录制消费端消息契约解析:从 recording_consumer Blueprint 到 ingest-replay-recordings 消费实现

Sentry Replays 录制消费端消息契约解析:从 recording_consumer Blueprint 到 ingest-replay-recordings 消费实现 Sentry Replays 录制消费端消息契约解析从 recording_consumer Blueprint 到 ingest-replay-recordings 消费实现【免费下载链接】sentryDeveloper-first error tracking and performance monitoring项目地址: https://gitcode.com/GitHub_Trending/sen/sentrySentry 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.mdSDK 事件与录制分段事件契约与 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 为什么用 msgpackBlueprint 开篇强调所有消息均以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 RecordingsReplay Recording Not Chunked3.1 字段契约小录制只有一种消息类型replay_recording_not_chunked其字段定义如下来自 recording_consumer.md字段类型说明typestring字面量replay_recording_not_chunkedreplay_idstring回放会话的唯一 IDkey_idOptional[int]上报所用的 API Key ID可选org_idint组织 IDproject_idint项目 IDreceivedint事件被接收时的 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 RecordingsChunk 与汇总消息大录制场景下共两种消息类型配合完成分块上传 → 汇总声明两步。4.1 Replay Recording Chunk录制分块字段类型说明typestring字面量replay_recording_chunkreplay_idstring回放会话的唯一 IDproject_idint项目 IDchunk_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_id、org_id、received字段也没有JSON 头部前缀payload 直接是分块字节id字段用于把多个 chunk 关联到同一条录制汇总消息上。4.2 Replay Recording录制汇总字段类型说明typestring字面量replay_recordingreplay_idstring回放会话的唯一 IDkey_idOptional[int]上报所用的 API Key ID可选org_idint组织 IDproject_idint项目 IDreceivedint事件被接收时的 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_recordinginstrumented_task( namePROCESS_REPLAY_RECORDING_TASK_NAME, namespacereplays_raw_tasks, retryRetry(times3, delay5), silo_modeSiloMode.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_iddecompress_segmentzlib 解压含未压缩兜底提取可选的replay_eventreplay 事件 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 前缀即为 TTLStorageBlob注释说明支持 30/60/90 天三种保留期需在存储桶上启用 TTL。读取侧 API 也遵循同一命名规则例如 api.md 中的 recording-segments 查询接口。六、打包格式录制与视频的二进制容器当消息携带replay_video时消费端会将 rrweb 录制与视频合包压缩后落库见 pack_replay_video。二进制容器格式定义在 pack.py首字节为 8 位类型标记0表示 rrweb1表示 video为兼容历史数据类型字节最大值限定为 9091 是 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_chunkreplay_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),仅供参考
返回列表