Neon WAL 处理全链路解析从 Compute 到 Pageserver 的深度解码与存储摘要本文深入剖析 Neon 数据库架构中 WALWrite-Ahead Log从 Compute 节点产生经 Safekeeper 接收、解码、转发最终被 Pageserver 摄入并存储的完整流程。重点解析 Safekeeper 的流解码与深度解码两层机制以及 Pageserver 如何利用 InterpretedWalRecord 实现计算存储分离避免传统 WAL Redo 的架构设计。阶段一Compute → SafekeeperWAL PushCompute 节点 (PostgreSQL)│ TCP连接, 发送 START_WAL_PUSH 命令▼1. 网络接入层wal_service.rs:61handle_socket()→ 创建SafekeeperPostgresHandler启动 PG 协议主循环handler.rs:306process_query()→ 匹配START_WAL_PUSH→ 调用handle_start_wal_push()2. WAL 流读取receive_wal.rs:191handle_start_wal_push_guts()建立 mpsc channelmsg_tx/msg_rx将网络读取和 WAL 处理解耦启动NetworkReader读网络消息 → 发到 channel启动WalAcceptor::run()从 channel 收消息 → 处理receive_wal.rs:412read_network_loop()→ 循环读取CopyData消息 →ProposerAcceptorMessage::parse()反序列化 →msg_tx.send()3. WAL 写入磁盘 流解码receive_wal.rs:600WalAcceptor::run()→ 收到AppendRequest→tli.process_msg(NoFlushAppendRequest)timeline.rs:1106→safekeeper.rs:937SafeKeeper::process_msg()→safekeeper.rs:1293handle_append_request()wal_storage.rs:437PhysicalStorage::write_wal()—— 这是第一个关键点self.write_exact(startpos, buf)?; // ① 写入磁盘 self.decoder.feed_bytes(buf); // ② 喂入解码器 while let Some((lsn, _rec)) self.decoder.poll_decode()? { // ③ 找记录边界 self.write_record_lsn lsn; }这里只做WAL page header 校验 记录边界识别 CRC 校验不做内容解析。对应lib.rs:409feed_bytes()→ 追加字节到内部 bufferlib.rs:413poll_decode()→ 版本分派到waldecoder_handler.rs:100poll_decode_internal()验证 page headermagic、LSN、continuation record 标志状态机WaitingForRecord→ 读xl_tot_len→ 跨页重组 →ReassemblingRecordcomplete_record()→ CRC 校验 → 处理XLOG_SWITCH此时 safekeeper 得到的是原始 WAL 字节Bytes只确认了记录边界未解析内容。阶段二Safekeeper → PageserverWAL 转发 深度解码Pageserver (WAL consumer)│ TCP连接, 发送 START_REPLICATION▼4. Safekeeper 侧解读 WAL 并发送send_wal.rs:443handle_start_replication()→send_wal.rs:559解读模式// send_wal.rs:601 let reader_handle InterpretedWalReader::spawn( wal_stream, start_pos, tx, shard, pg_version, ... ); // send_wal.rs:634 let sender_handle InterpretedWalSender::spawn(rx, pgb_writer, ...);send_interpreted_wal.rs:254InterpretedWalReader::spawn()→send_interpreted_wal.rs:393run_impl()let mut wal_decoder WalStreamDecoder::new(start_pos, self.pg_version); loop { let wal self.wal_stream.read().await?; // 从磁盘读 WAL wal_decoder.feed_bytes(wal); // 喂入流解码器 while let Ok(Some((lsn, recdata))) wal_decoder.poll_decode() { // ← 你选中的行 // 这里是深度解码 let records InterpretedWalRecord::from_bytes_filtered( recdata, shard_ids, next_record_lsn, self.pg_version )?; tx.send(Batch { records, ... }).await?; // 发送给 sender } }5. 深度解码decode_wal_record()核心decoder.rs:23InterpretedWalRecord::from_bytes_filtered()// ① 调用 postgres_ffi 的 decode_wal_record 解析 WAL 记录内容 decode_wal_record(buf, mut decoded, pg_version)?; // ② 按 shard 过滤并提取元数据 页面数据 MetadataRecord::from_decoded_filtered(decoded, pg_version)? // 元数据(事务/DDL等) SerializedValueBatch::from_decoded_filtered(decoded, ...) // 页面KV数据walrecord.rs:246decode_wal_record()—— 深度解析解析 XLogRecord headerrmid, info, xid遍历 block headers提取 RelFileNode、forknum、blkno、full-page image、压缩信息处理XLR_BLOCK_ID_ORIGIN、XLR_BLOCK_ID_TOPLEVEL_XID等特殊标志decoder.rs:77MetadataRecord::from_decoded_filtered()→ 按xl_rmid分发rmgr ID处理函数产物RM_XACT_IDdecode_xact_record()XACT commit/abort 信息RM_HEAP_IDdecode_heapam_record()VM bit 清除RM_DBASE_IDdecode_dbase_record()DB create/dropRM_SMGR_IDdecode_smgr_record()SMGR create/truncateRM_CLOG_IDdecode_clog_record()CLOG zero/truncateRM_LOGICALMSG_IDdecode_logical_message_record()LogicalMessage(Put/Failpoint)6. 网络发送给 Pageserversend_interpreted_wal.rs:650InterpretedWalSender::run()→ 从 channel 收Batch→ 序列化为InterpretedWalRecordswire format → 通过 PG 协议发送阶段三Pageserver 接收 存储Safekeeper 发送 InterpretedWalRecords (wire format)│ TCP▼7. Pageserver 接收 WALwalreceiver_connection.rs:116handle_walreceiver_connection()连接 safekeeper物理复制协议IDENTIFY_SYSTEM获取 WAL 位点START_REPLICATION PHYSICAL startpoint开始流复制创建WalIngest实例walreceiver_connection.rs:337接收循环ReplicationMessage::RawInterpretedWalRecords(raw) { let batch InterpretedWalRecords::from_wire(raw.data(), format, compression)?; for record in batch.records { walingest.ingest_record(interpreted, mut modification, ctx)?; } modification.commit(ctx).await?; }8. WAL 摄入引擎walingest.rs:233WalIngest::ingest_record()—— 中央分发ingest_record() ├── MetadataRecord 分发 │ ├── Xact(commit/abort) → ingest_xact_record() → 写 CLOG、删关系 │ ├── Dbase(create) → ingest_xlog_dbase_create() → 复制源库关系 │ ├── Dbase(drop) → ingest_xlog_dbase_drop() → 删除库关系 │ ├── Smgr(create) → ingest_xlog_smgr_create() → 创建关系元数据 │ ├── Smgr(truncate) → ingest_xlog_smgr_truncate() │ ├── Clog(truncate) → ingest_clog_truncate() │ ├── Heap(ClearVM) → ingest_clear_vm_bits() → WAL record 写入 │ ├── LogicalMessage(put) → ingest_logical_message_put() │ └── ... │ ├── modification.ingest_batch(interpreted.batch) → 页面数据累积 │ └── [pgdatadir_mapping.rs:1929] → pending_data_batch 累积 KV │ └── modification.put_checkpoint() → 更新 checkpoint (nextXid等)9. 提交到存储层pgdatadir_mapping.rs:2867DatadirModification::commit()let mut writer self.tline.writer().await; // 获取 TimelineWriter writer.put_batch(batch, ctx)?; // 写入 InMemoryLayer writer.finish_write(pending_lsn)?; // 推进 last_record_lsntimeline.rs:7862→inmemory_layer.rs:571InMemoryLayer::put_batch()file.write_raw(raw)→ 追加序列化字节到文件index.write()→ 更新内存 BTree 索引key → LSN/offset/length完整调用栈总览┌─────────────────────────────────────────────────────────────────────────┐ │ Compute 节点 │ │ PostgreSQL WAL writer → TCP → START_WAL_PUSH │ └────────────────────────────┬────────────────────────────────────────────┘ ▼ ┌─ Safekeeper ───────────────────────────────────────────────────────────┐ │ wal_service.rs:61 handle_socket() TCP accept │ │ handler.rs:306 process_query() PG协议分发 │ │ receive_wal.rs:263 read_first_message() 读首条消息 │ │ receive_wal.rs:412 read_network_loop() CopyData循环读取 │ │ receive_wal.rs:600 WalAcceptor::run() 消息处理循环 │ │ safekeeper.rs:1293 handle_append_request() WAL追加处理 │ │ wal_storage.rs:437 write_wal() ─────────────────────┐ │ │ ├── write_exact() 写入磁盘 │ │ │ ├── feed_bytes() 喂入WalStreamDecoder │ ← 阶段一 │ │ └── poll_decode() 找记录边界CRC校验 ← 你选中的行 │ (写入) │ │ (waldecoder_handler.rs:100) │ │ │ │ │ │ ─────────────── safekeeper 两条路径的分界线 ──────────── │ │ │ │ │ │ send_wal.rs:443 handle_start_replication() │ │ │ send_interpreted_wal:254 InterpretedWalReader::spawn()│ │ │ send_interpreted_wal:393 run_impl() ─────────────────────┐ │ │ │ ├── feed_bytes() 喂入WalStreamDecoder │ │ │ │ ├── poll_decode() 找记录边界 ← 又是你选中的行 │ │ ← 阶段二 │ │ ├── from_bytes_filtered() ★ 深度解码入口 │ │ (转发) │ │ │ ├── decode_wal_record() 解析WAL内容(walrecord:246)│ │ │ │ │ ├── MetadataRecord解码 (decoder.rs:77) │ │ │ │ │ └── SerializedValueBatch 页面数据序列化 │ │ │ │ └── tx.send(Batch) 发送给InterpretedWalSender │ │ │ │ │ │ │ send_interpreted_wal:650 InterpretedWalSender::run()│ │ │ └── BeMessage::InterpretedWalRecords 序列化网络发送│ │ └────────────────────────────┬────────────────────────────────────────────┘ │ TCP (InterpretedWalRecords wire format) ▼ ┌─ Pageserver ───────────────────────────────────────────────────────────┐ │ walreceiver_connection:116 handle_walreceiver_connection() 连接握手 │ │ walreceiver_connection:346 from_wire() 反序列化 │ │ walreceiver_connection:446 ingest_record() ───────────────┐ │ │ │ │ │ walingest.rs:233 WalIngest::ingest_record() 中央分发 │ │ │ ├── ingest_xact_record() → CLOG写入 关系删除 │ │ │ ├── ingest_xlog_dbase_create() → 复制数据库关系 │ ← 阶段三 │ │ ├── ingest_xlog_smgr_create() → 创建关系元数据 │ (存储) │ │ ├── ingest_clear_vm_bits() → VM页面WAL记录 │
Neon wal日志处理流程
Neon WAL 处理全链路解析从 Compute 到 Pageserver 的深度解码与存储摘要本文深入剖析 Neon 数据库架构中 WALWrite-Ahead Log从 Compute 节点产生经 Safekeeper 接收、解码、转发最终被 Pageserver 摄入并存储的完整流程。重点解析 Safekeeper 的流解码与深度解码两层机制以及 Pageserver 如何利用 InterpretedWalRecord 实现计算存储分离避免传统 WAL Redo 的架构设计。阶段一Compute → SafekeeperWAL PushCompute 节点 (PostgreSQL)│ TCP连接, 发送 START_WAL_PUSH 命令▼1. 网络接入层wal_service.rs:61handle_socket()→ 创建SafekeeperPostgresHandler启动 PG 协议主循环handler.rs:306process_query()→ 匹配START_WAL_PUSH→ 调用handle_start_wal_push()2. WAL 流读取receive_wal.rs:191handle_start_wal_push_guts()建立 mpsc channelmsg_tx/msg_rx将网络读取和 WAL 处理解耦启动NetworkReader读网络消息 → 发到 channel启动WalAcceptor::run()从 channel 收消息 → 处理receive_wal.rs:412read_network_loop()→ 循环读取CopyData消息 →ProposerAcceptorMessage::parse()反序列化 →msg_tx.send()3. WAL 写入磁盘 流解码receive_wal.rs:600WalAcceptor::run()→ 收到AppendRequest→tli.process_msg(NoFlushAppendRequest)timeline.rs:1106→safekeeper.rs:937SafeKeeper::process_msg()→safekeeper.rs:1293handle_append_request()wal_storage.rs:437PhysicalStorage::write_wal()—— 这是第一个关键点self.write_exact(startpos, buf)?; // ① 写入磁盘 self.decoder.feed_bytes(buf); // ② 喂入解码器 while let Some((lsn, _rec)) self.decoder.poll_decode()? { // ③ 找记录边界 self.write_record_lsn lsn; }这里只做WAL page header 校验 记录边界识别 CRC 校验不做内容解析。对应lib.rs:409feed_bytes()→ 追加字节到内部 bufferlib.rs:413poll_decode()→ 版本分派到waldecoder_handler.rs:100poll_decode_internal()验证 page headermagic、LSN、continuation record 标志状态机WaitingForRecord→ 读xl_tot_len→ 跨页重组 →ReassemblingRecordcomplete_record()→ CRC 校验 → 处理XLOG_SWITCH此时 safekeeper 得到的是原始 WAL 字节Bytes只确认了记录边界未解析内容。阶段二Safekeeper → PageserverWAL 转发 深度解码Pageserver (WAL consumer)│ TCP连接, 发送 START_REPLICATION▼4. Safekeeper 侧解读 WAL 并发送send_wal.rs:443handle_start_replication()→send_wal.rs:559解读模式// send_wal.rs:601 let reader_handle InterpretedWalReader::spawn( wal_stream, start_pos, tx, shard, pg_version, ... ); // send_wal.rs:634 let sender_handle InterpretedWalSender::spawn(rx, pgb_writer, ...);send_interpreted_wal.rs:254InterpretedWalReader::spawn()→send_interpreted_wal.rs:393run_impl()let mut wal_decoder WalStreamDecoder::new(start_pos, self.pg_version); loop { let wal self.wal_stream.read().await?; // 从磁盘读 WAL wal_decoder.feed_bytes(wal); // 喂入流解码器 while let Ok(Some((lsn, recdata))) wal_decoder.poll_decode() { // ← 你选中的行 // 这里是深度解码 let records InterpretedWalRecord::from_bytes_filtered( recdata, shard_ids, next_record_lsn, self.pg_version )?; tx.send(Batch { records, ... }).await?; // 发送给 sender } }5. 深度解码decode_wal_record()核心decoder.rs:23InterpretedWalRecord::from_bytes_filtered()// ① 调用 postgres_ffi 的 decode_wal_record 解析 WAL 记录内容 decode_wal_record(buf, mut decoded, pg_version)?; // ② 按 shard 过滤并提取元数据 页面数据 MetadataRecord::from_decoded_filtered(decoded, pg_version)? // 元数据(事务/DDL等) SerializedValueBatch::from_decoded_filtered(decoded, ...) // 页面KV数据walrecord.rs:246decode_wal_record()—— 深度解析解析 XLogRecord headerrmid, info, xid遍历 block headers提取 RelFileNode、forknum、blkno、full-page image、压缩信息处理XLR_BLOCK_ID_ORIGIN、XLR_BLOCK_ID_TOPLEVEL_XID等特殊标志decoder.rs:77MetadataRecord::from_decoded_filtered()→ 按xl_rmid分发rmgr ID处理函数产物RM_XACT_IDdecode_xact_record()XACT commit/abort 信息RM_HEAP_IDdecode_heapam_record()VM bit 清除RM_DBASE_IDdecode_dbase_record()DB create/dropRM_SMGR_IDdecode_smgr_record()SMGR create/truncateRM_CLOG_IDdecode_clog_record()CLOG zero/truncateRM_LOGICALMSG_IDdecode_logical_message_record()LogicalMessage(Put/Failpoint)6. 网络发送给 Pageserversend_interpreted_wal.rs:650InterpretedWalSender::run()→ 从 channel 收Batch→ 序列化为InterpretedWalRecordswire format → 通过 PG 协议发送阶段三Pageserver 接收 存储Safekeeper 发送 InterpretedWalRecords (wire format)│ TCP▼7. Pageserver 接收 WALwalreceiver_connection.rs:116handle_walreceiver_connection()连接 safekeeper物理复制协议IDENTIFY_SYSTEM获取 WAL 位点START_REPLICATION PHYSICAL startpoint开始流复制创建WalIngest实例walreceiver_connection.rs:337接收循环ReplicationMessage::RawInterpretedWalRecords(raw) { let batch InterpretedWalRecords::from_wire(raw.data(), format, compression)?; for record in batch.records { walingest.ingest_record(interpreted, mut modification, ctx)?; } modification.commit(ctx).await?; }8. WAL 摄入引擎walingest.rs:233WalIngest::ingest_record()—— 中央分发ingest_record() ├── MetadataRecord 分发 │ ├── Xact(commit/abort) → ingest_xact_record() → 写 CLOG、删关系 │ ├── Dbase(create) → ingest_xlog_dbase_create() → 复制源库关系 │ ├── Dbase(drop) → ingest_xlog_dbase_drop() → 删除库关系 │ ├── Smgr(create) → ingest_xlog_smgr_create() → 创建关系元数据 │ ├── Smgr(truncate) → ingest_xlog_smgr_truncate() │ ├── Clog(truncate) → ingest_clog_truncate() │ ├── Heap(ClearVM) → ingest_clear_vm_bits() → WAL record 写入 │ ├── LogicalMessage(put) → ingest_logical_message_put() │ └── ... │ ├── modification.ingest_batch(interpreted.batch) → 页面数据累积 │ └── [pgdatadir_mapping.rs:1929] → pending_data_batch 累积 KV │ └── modification.put_checkpoint() → 更新 checkpoint (nextXid等)9. 提交到存储层pgdatadir_mapping.rs:2867DatadirModification::commit()let mut writer self.tline.writer().await; // 获取 TimelineWriter writer.put_batch(batch, ctx)?; // 写入 InMemoryLayer writer.finish_write(pending_lsn)?; // 推进 last_record_lsntimeline.rs:7862→inmemory_layer.rs:571InMemoryLayer::put_batch()file.write_raw(raw)→ 追加序列化字节到文件index.write()→ 更新内存 BTree 索引key → LSN/offset/length完整调用栈总览┌─────────────────────────────────────────────────────────────────────────┐ │ Compute 节点 │ │ PostgreSQL WAL writer → TCP → START_WAL_PUSH │ └────────────────────────────┬────────────────────────────────────────────┘ ▼ ┌─ Safekeeper ───────────────────────────────────────────────────────────┐ │ wal_service.rs:61 handle_socket() TCP accept │ │ handler.rs:306 process_query() PG协议分发 │ │ receive_wal.rs:263 read_first_message() 读首条消息 │ │ receive_wal.rs:412 read_network_loop() CopyData循环读取 │ │ receive_wal.rs:600 WalAcceptor::run() 消息处理循环 │ │ safekeeper.rs:1293 handle_append_request() WAL追加处理 │ │ wal_storage.rs:437 write_wal() ─────────────────────┐ │ │ ├── write_exact() 写入磁盘 │ │ │ ├── feed_bytes() 喂入WalStreamDecoder │ ← 阶段一 │ │ └── poll_decode() 找记录边界CRC校验 ← 你选中的行 │ (写入) │ │ (waldecoder_handler.rs:100) │ │ │ │ │ │ ─────────────── safekeeper 两条路径的分界线 ──────────── │ │ │ │ │ │ send_wal.rs:443 handle_start_replication() │ │ │ send_interpreted_wal:254 InterpretedWalReader::spawn()│ │ │ send_interpreted_wal:393 run_impl() ─────────────────────┐ │ │ │ ├── feed_bytes() 喂入WalStreamDecoder │ │ │ │ ├── poll_decode() 找记录边界 ← 又是你选中的行 │ │ ← 阶段二 │ │ ├── from_bytes_filtered() ★ 深度解码入口 │ │ (转发) │ │ │ ├── decode_wal_record() 解析WAL内容(walrecord:246)│ │ │ │ │ ├── MetadataRecord解码 (decoder.rs:77) │ │ │ │ │ └── SerializedValueBatch 页面数据序列化 │ │ │ │ └── tx.send(Batch) 发送给InterpretedWalSender │ │ │ │ │ │ │ send_interpreted_wal:650 InterpretedWalSender::run()│ │ │ └── BeMessage::InterpretedWalRecords 序列化网络发送│ │ └────────────────────────────┬────────────────────────────────────────────┘ │ TCP (InterpretedWalRecords wire format) ▼ ┌─ Pageserver ───────────────────────────────────────────────────────────┐ │ walreceiver_connection:116 handle_walreceiver_connection() 连接握手 │ │ walreceiver_connection:346 from_wire() 反序列化 │ │ walreceiver_connection:446 ingest_record() ───────────────┐ │ │ │ │ │ walingest.rs:233 WalIngest::ingest_record() 中央分发 │ │ │ ├── ingest_xact_record() → CLOG写入 关系删除 │ │ │ ├── ingest_xlog_dbase_create() → 复制数据库关系 │ ← 阶段三 │ │ ├── ingest_xlog_smgr_create() → 创建关系元数据 │ (存储) │ │ ├── ingest_clear_vm_bits() → VM页面WAL记录 │