设计解析:从两遍式 Segment/Checkpoint 到源码级恢复机制)
Loki 写入预写日志WAL设计解析从两遍式 Segment/Checkpoint 到源码级恢复机制【免费下载链接】lokiLike Prometheus, but for logs.项目地址: https://gitcode.com/GitHub_Trending/lok/loki导读本文以 Loki 官方设计文档 2020-09-Write-Ahead-Log.md 为主体系统讲解 Loki 如何通过写入预写日志Write-Ahead LogWAL为 ingester 组件提供本地磁盘级的写入持久化保障先用segments段记录已接受的写入再用checkpoints检查点把段记录合并为更高效的表示以加速回放。读完本文你将掌握 WAL 的段文件与检查点生命周期、记录类型的编码细节、故障后的恢复replay流程以及对应的全部配置参数与生产部署要点并能对照 pkg/ingester/wal.go、pkg/ingester/checkpoint.go、pkg/ingester/recovery.go 等源码理解其底层实现。背景与动机ImpetusLoki 在写入路径上早已采取多种措施保证日志数据的持久性最典型的是在 ingester 中配置可调的复制因子replication factor来提供冗余。但对于持久性保证而言这仍然留有缺口——尤其是单二进制single binary部署场景一旦 ingester 进程崩溃或重启尚未刷入长期存储的内存数据就可能丢失。该设计文档提出的方案是在 ingester 组件上通过本地磁盘引入一个写入预写日志WAL用先落盘、再进内存的方式为已接受的写入提供一条可回放replay的持久化通道从而对既有复制机制形成补充。总体策略两遍式 WAL设计文档建议实现一个两遍式two passWAL第一遍segments段。逐条记录已被接受的写入是最基础的 WAL。仅凭段文件即可在无任何外部输入的情况下重建 ingester 的内存状态。第二遍checkpoints检查点。把第一遍的记录合并成更高效的表示从而大幅加快回放速度。这两遍配合定期截断机制共同解决WAL 无限增长与回放太慢两个问题。Segments段文件段是第一遍也是最基础的 WAL。每个段存储若干条已被接受的写入记录。关键特性如下段大小每个段是 32kB 的若干倍一个段写满后会自动创建下一个段段在磁盘上按序号依次命名。设计初稿提出先尝试 256KB 的段大小再按需调整。实际实现在 pkg/ingester/wal.go#L27 中段大小被定义为wlog.DefaultSegmentSize * 4即基于 Prometheustsdb/wlog包默认段大小的 4 倍——这正体现了尽量复用 Prometheus WAL 包的实现目标。目录布局段写入在 WAL 目录下形如data └── wal ├── 000000 ├── 000001 └── 000002写入顺序从 pkg/ingester/wal.go#L118-L151 的Log方法可以看到每个记录总是先写 Series序列再写 Entries日志条目且写入时复用共享的字节缓冲池recordPool写完后归还。Truncation定期截断为防止 WAL 无限增长同时把**已刷入存储flushed**的操作从 WAL 中移除WAL 会按可配置间隔ingester.checkpoint-duration定期截断删除除当前活跃段之外的所有段。这正是检查点登场的时机——截断前必须确保当前正在写入的段不被误删。Checkpoints检查点机制检查点生成流程在截断 WAL 之前系统会先把 WAL 段向前推进一个NextSegment以确保不会删除正在写入的段。此时目录形如data └── wal ├── 000000 ├── 000001 ├── 000002 - 可能没写满没关系 └── 000003 - 新创建的、空的活跃段随后每个内存中的 stream 会在一个按摊销计算的时间间隔内被逐个遍历并写入检查点。该间隔的计算方式为checkpoint_duration / in_memory_streams在 pkg/ingester/checkpoint.go#L666 的源码中这一逻辑被实现为// Give a 10% buffer to the checkpoint duration in order to account for // new series, slow writes, etc. perSeriesDuration : (90 * c.dur) / (100 * time.Duration(n))即每个 stream 分配到的写入时间片为检查点时长的 90%预留 10% 缓冲应对新序列、慢写入等情况通过time.NewTicker(perSeriesDuration)把检查点写入在checkpoint_duration内均匀摊销避免瞬时 IOPS 尖峰。检查点写完后会从其临时目录被原子地移动到ingester.wal-dir并以它开始前最后一个段的序号命名checkpoint.000002随后删除所有已被覆盖的段000000、000001、000002以及任何更早的检查点。完成后目录形如data └── wal ├── checkpoint.000002 - 已完成的检查点 └── 000003 - 当前活跃的 wal 段源码中的原子性与清理细节从 pkg/ingester/checkpoint.go 可以印证整套流程的实现Advancecheckpoint.go#L319-L359先调用wlog.Segments找出最后一个段的序号调用segmentWAL.NextSegment()推进一个段确保不覆盖、也尽量少地在检查点之上重放段然后在临时目录checkpoint.XXXXXX.tmp中创建新的检查点 WAL同时会清理历史遗留的.tmp半成品检查点目录。Writecheckpoint.go#L364-L384每个 stream 序列化为一条带CheckpointRecord类型头的记录攒够约 1MB 批量刷盘bufSize 120。Closecheckpoint.go#L585-L621先把临时检查点目录通过fileutil.Replace原子重命名为目标目录checkpoint.000002再调用segmentWAL.Truncate(w.lastSegment 1)删除旧段并deleteCheckpoints(w.lastSegment)清理更早的检查点。若删除失败只记录错误、不阻塞后续流程留到下次再试。命名约定checkpointPrefix checkpoint.、checkpointTmpSuffix .tmpcheckpoint.go#L400-L401恢复时通过正则^checkpoint\.(\d)(\.tmp)?$识别并选取序号最大的已完成检查点。检查点操作的排队Queueing Checkpoint operations一个检查点操作有可能与另一个同时开始。设计文档明确指出此时正在运行的那个检查点操作应当忽略其内部 ticker尽快刷完所有 series然后下一个检查点操作才能开始。这会在摊销生效前造成局部的 IOPS 尖峰——这也是文档强调WAL 应当运行在独立磁盘上以缓解噪声邻居noisy neighbor问题的重要原因。检查点完成后旧检查点被回收reap。在源码中这个尽快刷完的切换点对应 checkpoint.go#L683-L690当time.Since(start) c.dur当前检查点已超过一个检查点时长的预算时immediate被置为true后续 series 不再等待 ticker直接连续写入。WAL 记录类型设计文档定义了两种核心 WAL 记录类型Streams序列记录当 ingester 收到一个内存中尚不存在的序列的 push 时写入一条Stream记录。高层结构如下type SeriesRecord struct { UserID string Labels labels.Labels Fingerprint uint64 // label fingerprint }Logs日志记录当 ingester 收到 push 时写入一条Logs记录包含其所属序列的 fingerprint 以及一组(timestamp, log_line)元组若适用该记录一定写在对应的Stream记录之后type LogsRecord struct { UserID string Fingperprint uint64 // label fingerprint for the series these logs refer to Entries []logproto.Entry }实际实现中的编码演进在 pkg/ingester/wal/encoding.go#L18-L31 中记录类型被实现为一个字节枚举且已演进出多个版本const ( WALRecordSeries RecordType iota 1 // 序列记录 WALRecordEntriesV1 // 日志条目 v1 CheckpointRecord // 检查点记录基于 protobuf WALRecordEntriesV2 // 日志条目 v2带 counter回放时可放宽排序约束 WALRecordEntriesV3 // 日志条目 v3带结构化元数据structured metadata )Loki 始终写入最新的WALRecordEntriesV3但向后兼容地读取所有旧版本CurrentEntriesRec。编码上的几个关键设计encoding.go#L99-L149时间戳差分记录中第一个条目的时间戳以 8 字节大端存储为基准其余条目只存与基准的差值PutVarint64显著压缩空间变长整数条目数量、行长度、结构化元数据数量等全部使用Uvarint编码结构化元数据V3 起每条 entry 附带(name, value)键值对列表对应logproto.LabelAdapter流式解析DecodeRecord依据类型头分派到DecodeEntries/series 解码器未知类型直接报错。恢复Restoration回放 WAL回放流程回放 WAL 的方式是先把可用的检查点载入内存再在其上依次重放段文件中的操作checkpoint.000003→000004→000005…。由于摊销引入的延迟部分操作可能已在检查点中重复出现重放时会失败但这没关系——不会丢失任何数据只是某些数据被写了两次而重复写入会被忽略。源码级的多路并行恢复恢复逻辑在 pkg/ingester/recovery.go 中实现与设计文档完全对应先检查点、后段newCheckpointReaderrecovery.go#L36-L55通过lastCheckpoint找到序号最大的已完成检查点若不存在则用NoopWALReader跳过。启动恢复前还会主动清理遗留的.tmp半成品检查点和已被取代的旧检查点cleanupCheckpointsAtStartup。两类恢复入口RecoverCheckpointrecovery.go#L307-L353把检查点记录反序列化为Series含 chunks并重建RecoverWALrecovery.go#L243-L305先SetStream处理所有序列确保不把条目写到不存在的序列上再按 fingerprint 哈希分发 entries。并行 workerrecoverGenericrecovery.go#L360-L427根据runtime.GOMAXPROCS(0)启动多个 worker按fingerprint % nWorkers哈希把恢复输入分发到各 worker 通道实现多核并行回放。重复写入容错Push时recovery.go#L202-L208对ErrEntriesExist这类乱序/重复错误直接忽略仅累加duplicateEntriesTotal指标——这正是文档所说写两次会被忽略的落地实现。内存上限回放期间还受ingester.wal-replay-memory-ceiling默认 4GBpkg/ingester/wal.go#L28约束超过该上限会先向存储刷数据再继续回放避免恢复过程把内存打爆。配置参数与部署要点配置参数WAL 的全部配置集中在 pkg/ingester/wal.go#L30-L60 的WALConfig与RegisterFlags中同时支持 YAML 与命令行 flag配置项YAML命令行 flag默认值说明enabledingester.wal-enabledtrue是否启用 WAL 写入禁用时使用noopWAL空实现diringester.wal-dirwalWAL 数据存储与恢复目录checkpoint_durationingester.checkpoint-duration5m检查点创建间隔启用 WAL 时该值必须 ≥ 1flush_on_shutdowningester.flush-on-shutdownfalseWAL 启用时关闭进程是否把 chunks 刷入长期存储replay_memory_ceilingingester.wal-replay-memory-ceiling4GB回放期间允许 WAL 使用的最大内存超限则先刷盘再继续支持 KB/MB/GB 后缀disk_full_thresholdingester.wal-disk-full-threshold0.90磁盘使用率阈值0.0~1.0达到后 WAL 节流throttle写入设为 0 关闭节流其中disk_full_threshold与replay_memory_ceiling属于实现中新增的防护能力monitorDiskpkg/ingester/wal.go#L194-L239每 10 秒检查一次磁盘用量超过阈值即置位diskThrottled配合IsDiskThrottled()让写入路径降速防止磁盘写满拖垮整个 ingester。部署要求设计文档对 WAL 的部署提出两条硬性要求持久化磁盘引入 WAL 要求 ingester 拥有跨重启保持连接的持久磁盘——这在 Kubernetes 中与StatefulSet是天然契合的组合独立磁盘强烈建议 WAL 使用独立磁盘使其既不受其他组件影响、也不会成为噪声邻居尤其是在任何 IOPS 尖峰期间例如多个检查点排队时的瞬时刷盘。实现目标与方案取舍实现目标Implementation goals设计文档明确了三条实现目标且均已落地复用 Prometheus WAL 包在可行时直接使用github.com/prometheus/prometheus/tsdb/wlog接口负责页对齐page alignment并使用[]byte同时确保该包能处理任意长的记录Loki 场景下的超长日志行。这一点在 pkg/ingester/wal.go#L14 的 import 中得到印证。内存表示 ↔[]byte高效互转保证内存中的序列/chunk 能高效地序列化/反序列化以便从检查点快速加载——对应toWireChunks/fromWireChunks与SerializeForCheckpointTocheckpoint.go#L46-L112。已刷盘 chunk 的保留即使经过 WAL 回放已刷入存储的 chunk 也要保留ingester.retain-period时长这一约束由下游 flush 路径配合实现。方案取舍Alternatives设计文档比较了两种备选方案并说明为何没有采纳方案一直接使用 Cortex WAL。由于本方案不是从 WAL 记录重建检查点而是对内存做 dump因此瓶颈不在吞吐而在内存大小可以仅按时长触发检查点而无需额外核算吞吐与 Cortex WAL 方案几乎等价。唯一的隐患是两个检查点之间会积累大量段数据日志吞吐波动大若按时长检查点不够用未来可能再考虑其他手段。方案二不从内存构建检查点而是新增Blocks与FlushedChunks两种 WAL 记录类型。前者在压缩块切分后写入整块后者在 chunk 刷盘时写入整个 chunk 及其包含的块序列。这能把写入均匀摊销块切分与 chunk 刷盘都被假设为均匀分布且用 jitter 同步还能利用ingester.retain-period丢弃已过期的 WAL 记录以加速回放type FlushRecord struct { Fingerprint uint64 // labels FlushedAt uint64 // timestamp when it was flushed, can be used with ingester.retain-period to either keep or discard records on replay LastEntry logproto.Entry // last entry included in the flushed chunk }该方案还能在不依赖 ingester 内部状态的情况下构建检查点但代价是需要按记录类型分区的多个 WAL以便按FlushedChunks → Blocks → Series → Samples的优先级迭代、对低优先级类型做 no-op。文档的结论是收益抵不上成本——尤其是考虑到简单的替代方案已够用以及 ingester 内部结构变化时新增记录类型的扩展性成本。总结Loki 的 WAL 以段记录 检查点合并 定期截断 回放恢复为核心闭环为 ingester 提供了不依赖复制因子的本地持久化兜底段负责原样记录每次已接受的写入检查点按ingester.checkpoint-duration周期性地把内存状态快照到磁盘并回收旧段崩溃后先载入检查点、再按序号重放段文件即可完整还原内存状态。从 wal.go、checkpoint.go、recovery.go、wal/encoding.go 的实现来看设计文档中的每一项策略都有对应落地Prometheustsdb/wlog的复用、.tmp目录的原子重命名、检查点的摊销写入与超时立即模式、多 worker 并行恢复、重复写入容忍以及磁盘满节流、回放内存上限等现代防护手段。生产部署时请务必为 ingester 配置持久化且独立的磁盘Kubernetes 下优先 StatefulSet并按业务吞吐合理调整ingester.checkpoint-duration与ingester.wal-dir即可在单二进制与微服务两种形态下获得一致的数据持久性保障。【免费下载链接】lokiLike Prometheus, but for logs.项目地址: https://gitcode.com/GitHub_Trending/lok/loki创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考