Documentation
¶
Overview ¶
Package wal 实现 SamKV 的顺序预写日志、记录校验和可配置持久性策略。
正常写入使用 PutRecord/DeleteRecord 构造记录,再交给 WalManger.AppendRecord。 SyncEveryWrite 在返回前执行文件 Sync,适合不能接受最近写入丢失的场景;SyncInterval 先写入内存缓冲,由周期任务或缓冲满触发 Sync,吞吐更高,但进程或主机崩溃时可能 丢失最近一个同步窗口。
单条编码记录可以大于 BufferSize:管理器会先刷出旧缓冲,再直接写入并 Sync,不会等待 一个永远无法容纳它的缓冲。Flush 可建立显式持久化边界,Replace 用于 Checkpoint 后 原子改写仍需恢复的记录,Close 会停止后台任务并尝试刷出剩余数据。
ReadRecord 校验长度、类型和 CRC32,最大接受 64 MiB payload。io.EOF 表示正常读到文件 末尾,io.ErrUnexpectedEOF 表示尾部记录不完整;ErrChecksum 表示完整长度内的数据损坏。 WalManger 可并发追加,但同一数据目录只能由上层 Store 的目录锁保证单进程独占。
Index ¶
- Constants
- Variables
- func ParseSegmentID(path string) (uint64, bool)
- func SegmentPath(dir string, id uint64) string
- type Options
- type Record
- type RecordType
- type RecoveryOptions
- type RecoveryReport
- type Segment
- type SyncPolicy
- type WalManger
- func (wm *WalManger) ActiveSegment() Segment
- func (wm *WalManger) AppendLog(data []byte) error
- func (wm *WalManger) AppendRecord(record *Record) error
- func (wm *WalManger) Background()
- func (wm *WalManger) Close() error
- func (wm *WalManger) Err() error
- func (wm *WalManger) Flush() error
- func (wm *WalManger) PruneThrough(maxID uint64) (int, error)
- func (wm *WalManger) Replace(data []byte) error
- func (wm *WalManger) Reset() error
- func (wm *WalManger) Rotate() (uint64, error)
- type WalWriter
Examples ¶
Constants ¶
const ( // DefaultSize 是周期同步模式的默认内存缓冲容量。 DefaultSize = 64 * 1024 // 64 KiB // FlushInterval 是周期同步模式的默认最大等待时间。 FlushInterval = 50 * time.Millisecond )
const ( // DefaultSegmentSize 是 WAL 单个 segment 的默认目标大小。 // 完整 record 不会跨 segment 拆分,因此文件可能被最后一次 append 略微撑大。 DefaultSegmentSize int64 = 64 << 20 )
Variables ¶
var ( // ErrInvalidRecord 表示记录类型、长度或字段组合不符合 WAL 格式。 ErrInvalidRecord = errors.New("invalid wal record") // ErrChecksum 表示 payload 长度完整但 CRC32 校验失败。 ErrChecksum = errors.New("wal checksum mismatch") // ErrRecordTooLarge 表示流中声明的 payload 超过 64 MiB 安全上限。 ErrRecordTooLarge = errors.New("wal record too large") )
var ErrAmbiguousLayout = errors.New("wal: legacy and segmented WAL files coexist")
ErrAmbiguousLayout 表示目录同时存在旧 wal.log 和新 segment,无法可靠判断追加顺序。
var ( // ErrCorruptSegment 表示 WAL segment 无法在已知 record 边界上继续恢复。 ErrCorruptSegment = errors.New("wal: corrupt segment") )
var ErrInvalidOptions = errors.New("wal: invalid options")
ErrInvalidOptions 表示缓冲大小、同步策略或同步间隔组合无效。
var ( // ErrInvalidPrune 表示调用方试图删除当前活动 segment 或更高编号。 ErrInvalidPrune = errors.New("wal: prune boundary reaches active segment") )
Functions ¶
func ParseSegmentID ¶
ParseSegmentID 从规范文件名中解析 segment ID。 path 可以是完整路径;编号必须为 20 位十进制且大于 0。
func SegmentPath ¶
SegmentPath 返回指定 ID 的规范 WAL segment 路径。 id=0 不是合法持久化编号,内部调用方必须从 1 开始分配。
Types ¶
type Options ¶
type Options struct {
// BufferSize 是周期模式的内存缓冲容量,必须大于 0;它不是单条记录大小上限。
BufferSize int
// SyncPolicy 决定 AppendLog/AppendRecord 返回前是否完成 fsync。
SyncPolicy SyncPolicy
// SyncInterval 仅供 SyncInterval 策略使用,必须大于 0;严格模式会忽略它。
SyncInterval time.Duration
// SegmentSize 是触发 WAL segment 轮转的目标字节数,必须大于 0。
SegmentSize int64
// SegmentMaxRecords 是单段最多容纳的 record 数;0 表示只按 SegmentSize 轮转。
SegmentMaxRecords uint64
}
Options 控制 WAL 的缓冲容量和持久性策略。
type Record ¶
type Record struct {
// Type 是写入或删除操作。
Type RecordType
// Sequence 由上层 Store 分配,用于恢复版本顺序;零值合法。
Sequence uint64
// Key 必须非空。
Key []byte
// Value 在 RecordPut 中可以为空,在 RecordDelete 中必须为空。
Value []byte
}
Record 是一条可校验的 WAL 操作。 PutRecord/DeleteRecord 不复制传入切片,调用方在 Encode 或 AppendRecord 返回前不得修改它们。
Example ¶
package main
import (
"fmt"
"github.com/23jdd/SamKv/pkg/wal"
)
func main() {
record := wal.PutRecord([]byte("key"), []byte("value"))
record.Sequence = 7
encoded, err := record.Encode()
if err != nil {
panic(err)
}
decoded, err := wal.Decode(encoded)
if err != nil {
panic(err)
}
fmt.Println(decoded.Type, decoded.Sequence, string(decoded.Key), string(decoded.Value))
}
Output: 1 7 key value
func ReadRecord ¶
ReadRecord 从流中读取并校验一条完整 WAL 记录。 干净 EOF 原样返回 io.EOF;截断帧返回 io.ErrUnexpectedEOF;payload 超过 64 MiB 返回 ErrRecordTooLarge。
type RecordType ¶
type RecordType uint8
RecordType 区分写入和删除记录;零值及未知值都不是有效磁盘类型。
const ( // RecordPut 表示为 Key 写入 Value。 RecordPut RecordType = iota + 1 // RecordDelete 表示删除 Key;该类型的 Value 必须为空。 RecordDelete )
type RecoveryOptions ¶
type RecoveryOptions struct {
// SkipCorruptedRecords 跳过长度完整但 checksum 或 record 内容无效的单条帧。
SkipCorruptedRecords bool
// RepairTrailingPartial 截断最后一个 segment 的半条尾记录;非末段永远不会自动截断。
RepairTrailingPartial bool
}
RecoveryOptions 控制 WAL 回放遇到损坏时的策略。
func DefaultRecoveryOptions ¶
func DefaultRecoveryOptions() RecoveryOptions
DefaultRecoveryOptions 返回适合 Store 启动恢复的容错策略。
type RecoveryReport ¶
type RecoveryReport struct {
Segments int
Records int
SkippedRecords int
TruncatedBytes int64
LastSegmentID uint64
}
RecoveryReport 汇总一次 segment replay 的可观测结果。
func ReplaySegments ¶
func ReplaySegments( dir string, options RecoveryOptions, apply func(*Record) error, ) (RecoveryReport, error)
ReplaySegments 按 ID 递增回放所有 WAL segment。 apply 只会收到校验和格式均有效的记录;apply 返回错误时立即停止且不修改磁盘。
type Segment ¶
Segment 描述一个已经发布到 WAL 目录的 segment。 ID 决定恢复顺序,Path 是绝对或基于传入 dir 拼出的路径,Size 是枚举时的文件大小快照。
func ListSegments ¶
ListSegments 返回目录中全部规范 WAL segment,并按 ID 严格递增排序。 目录不存在时返回 os.ErrNotExist;临时文件、旧 wal.log 和命名非法文件不会混入结果。
type SyncPolicy ¶
type SyncPolicy uint8
SyncPolicy 控制 AppendLog 返回前 WAL 数据需要达到的持久化程度。
const ( // SyncInterval 由后台任务按固定间隔执行 fsync。 // 写入返回时数据可能仍在操作系统页缓存中,崩溃时可能丢失最近一个同步周期的数据。 SyncInterval SyncPolicy = iota // SyncEveryWrite 要求 AppendLog 在返回前完成 fsync。 SyncEveryWrite )
type WalManger ¶
type WalManger struct {
Dir string
// contains filtered or unexported fields
}
WalManger 管理 WAL 的内存缓冲、顺序写入和后台刷盘。
Example ¶
package main
import (
"bytes"
"fmt"
"os"
"github.com/23jdd/SamKv/pkg/wal"
)
func main() {
dir, err := os.MkdirTemp("", "samkv-wal-example-")
if err != nil {
panic(err)
}
defer os.RemoveAll(dir)
options := wal.DefaultOptions()
options.SyncPolicy = wal.SyncEveryWrite
manager, err := wal.NewWithOptions(dir, options)
if err != nil {
panic(err)
}
if err := manager.AppendRecord(wal.PutRecord([]byte("durable"), []byte("yes"))); err != nil {
panic(err)
}
if err := manager.Close(); err != nil {
panic(err)
}
segments, err := wal.ListSegments(dir)
if err != nil || len(segments) != 1 {
panic("unexpected WAL segments")
}
data, err := os.ReadFile(segments[0].Path)
if err != nil {
panic(err)
}
record, err := wal.ReadRecord(bytes.NewReader(data))
if err != nil {
panic(err)
}
fmt.Println(string(record.Key), string(record.Value))
}
Output: durable yes
func NewWithOptions ¶
NewWithOptions 打开或创建最新 WAL segment,并按给定持久性策略启动后台任务。 dir 不存在时自动创建;只有旧 wal.log 时会迁移为 segment 1,同时存在新旧布局则拒绝打开。 成功后调用方必须调用 Close,即使后续 AppendRecord 返回错误。
func (*WalManger) ActiveSegment ¶
ActiveSegment 返回当前追加 segment 的快照。
func (*WalManger) AppendLog ¶
AppendLog 将已经编码的 WAL 数据追加到缓冲区。 data 可以为空;若希望日志可恢复,必须传入一条或多条完整 Record 编码。 大于配置缓冲容量的单条数据会先刷出旧缓冲,再直接写盘并 Sync,避免永久等待。 方法可并发调用;Close 开始后返回 os.ErrClosed,后台错误会在后续调用中返回。
func (*WalManger) AppendRecord ¶
AppendRecord 编码并追加一条 WAL 记录。 record=nil、空 key 或非法类型返回 ErrInvalidRecord;持久性由创建管理器时的 SyncPolicy 决定。
func (*WalManger) Background ¶
func (wm *WalManger) Background()
Background 定时把 WAL 缓冲同步到磁盘。 NewWithOptions 已自动启动它,外部不得再次调用,否则 Close 的 WaitGroup 计数不匹配。
func (*WalManger) Close ¶
Close 停止后台刷盘协程,刷出剩余 buffer,并关闭当前活动 segment。 Close 可重复调用并返回第一次关闭结果;开始关闭后新的追加返回 os.ErrClosed。
func (*WalManger) Flush ¶
Flush 将当前 WAL 内存 buffer 同步刷到活动 segment。 Checkpoint 前必须先 Flush,避免仍在内存中的 WAL 记录丢失。 空缓冲返回 nil;此前记录已经在对应写入路径完成 Sync。
func (*WalManger) PruneThrough ¶
PruneThrough 删除 ID <= maxID 的已封存 segment。 调用方必须先持久化引用这些记录的 Manifest;重复调用会忽略已经不存在的旧段。
func (*WalManger) Replace ¶
Replace 把 data 写入更高 ID 的新 segment,Sync 并 Rename 发布后才切换活动句柄和删除旧段。 data 必须是零条或多条完整编码 record;崩溃发生在发布前会保留旧段,发布后但回收前会保留新旧两组。 Store 的新 Checkpoint 路径使用 Rotate+PruneThrough;Replace 仅保留给兼容调用方。