wal
# 1.简介
# 1.1.WAL日志
WAL全称为write ahead log预写日志,类似于redo log,是etcd中用于确保数据持久化和恢复能力的关键机制。WAL的主要目的是数据变更被应用到持久存储之前,变更记录会记录到一个日志文件中,避免系统崩溃或意外中断带来的数据丢失问题,通过WAL日志重放恢复尚未持久化到数据库中的变更,保证数据的一致性和可靠性。type WAL struct { dir string // WAL文件保存的路径 dirFile *os.File // dir打开的一个目录fd对象 metadata []byte // 创WAL时传入的字节序列,主要是节点ID和集群ID信息,每创建WAL文件就会写到首部 state raftpb.HardState // append过程中保存的硬状态信息,WAL有切割时会在新的WAL头部保存最新的 start walpb.Snapshot // 记录最后一次保存的snapshot信息,主要是snapshot中末尾日志的index和term decoder *decoder // 读取WAL日志文件时,将protobuf反序列化为Record实例 readClose func() error // 用于关闭decoder关联的reader,关闭WAL读模式 enti uint64 // 最后一次保存到WAL中的日志的index encoder *encoder // 写入WAL日志文件的Record实例序列化为protobuf locks []*fileutil.LockedFile // 当前WAL实例管理的所有WAL日志文件对应的句柄 fp *filePipeline // 负责创建新的临时文件 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
# 1.2.WAL文件
WAL内容最终会存在.wal文件,支持读取模式或追加模式。创建.wal文件时先在临时目录中初始化文件及写入元数据,通过锁定和预分配空间准备就绪,再将临时目录重命名为目标目录,同时同步父目录确保操作持久化。此外,文件都是先写再分割,切割文件时同样先在临时文件写入必要数据,最后再进行重命名。
补充
1.
wal文件大小限制为64MB,超出后会拆分2.旧文件被切分、临时文件重命名、服务器关闭、任期变化产生的新日志、保存快照都会涉及
WAL日志刷盘
# 1.3.Frame存储
WAL文件使用Frame格式组织日志条目,这种涉及既保证了数据完整性,又优化了读写性能。Frame的签名存储一个64bit的长度标识,根据8bit对齐。+---------------------+---------------------+-----------------------------+ | is_padding (1bit) | padding_len (7bit) | data_length (56bit) | +---------------------+---------------------+-----------------------------+ | 实际数据 (变长) | +-----------------------------------------------------------------------+ | CRC32 (4字节) | +-----------------------------------------------------------------------+1
2
3
4
5
6
7由于数据都是根据
8bit对齐,但是Encode的数据长度可能不是8的整数倍,因此为了补齐8的整数倍需要额外padding一定的长度,这部分长度不会超过7,因此可以按照3bit存,通过1bit标识是否需要padding,以便解码数据,decode时其实就是根据Frame格式解析出正确的data长度拿到内容。
# 2.日志存取
# 2.1.Record
更新
WAL数据时会向WAL日志写入metadata、snapshot等数据,其实都是Record为单位保存的,Record按照类型划分包含不同内容。// Record结构,WAL稳定存储的消息共两种,这是第一种普通存储格式 type Record struct { Type int64 `protobuf:"varint,1,opt,name=type" json:"type"` // 类型(以下4类) Crc uint32 `protobuf:"varint,2,opt,name=crc" json:"crc"` // 校验码 Data []byte `protobuf:"bytes,3,opt,name=data" json:"data,omitempty"` // 序列化后的数据 XXX_unrecognized []byte `json:"-"` // 保留字段,兼容未识别数据 } // Type取值包含以下几种 const ( metadataType int64 = iota + 1 entryType stateType crcType snapshotType // warnSyncDuration is the amount of time allotted to an fsync before // logging a warning warnSyncDuration = time.Second ) // 元数据 type Metadata struct { NodeID uint64 `protobuf:"varint,1,opt,name=NodeID" json:"NodeID"` // 节点ID ClusterID uint64 `protobuf:"varint,2,opt,name=ClusterID" json:"ClusterID"` // 集群ID XXX_unrecognized []byte `json:"-"` // 保留字段,兼容未识别数据 } // 日志 type Entry struct { Term uint64 `protobuf:"varint,2,opt,name=Term" json:"Term"` // 任期 Index uint64 `protobuf:"varint,3,opt,name=Index" json:"Index"` // 索引 Type EntryType `protobuf:"varint,1,opt,name=Type,enum=raftpb.EntryType" json:"Type"` // 类型 Data []byte `protobuf:"bytes,4,opt,name=Data" json:"Data,omitempty"` // 数据 XXX_unrecognized []byte `json:"-"` // 保留区 } // 硬状态 type HardState struct { Term uint64 `protobuf:"varint,1,opt,name=term" json:"term"` // 任期 Vote uint64 `protobuf:"varint,2,opt,name=vote" json:"vote"` // 索引 Commit uint64 `protobuf:"varint,3,opt,name=commit" json:"commit"` // 已提交日志索引 XXX_unrecognized []byte `json:"-"` // 保留区 } // 快照,这里的快照只是部分snapshot信息 type Snapshot struct { Index uint64 `protobuf:"varint,1,opt,name=index" json:"index"` // 末尾日志的index Term uint64 `protobuf:"varint,2,opt,name=term" json:"term"` // 末尾日志的term XXX_unrecognized []byte `json:"-"` // 保留区 }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51一条
Record需要先序列化后才能持久化,通过encode函数完成。func (e *encoder) encode(rec *walpb.Record) error { e.mu.Lock() defer e.mu.Unlock() // 计算crc写入record的crd字段 e.crc.Write(rec.Data) rec.Crc = e.crc.Sum32() var ( data []byte err error n int ) // 超出预分配的buf时,使用动态分配 if rec.Size() > len(e.buf) { data, err = rec.Marshal() } else { // 否则使用预分配的buf n, err = rec.MarshalTo(e.buf) data = e.buf[:n] } // 帧对齐处理,涉及数据长度和填充信息编码 lenField, padBytes := encodeFrameSize(len(data)) // 先写recode编码后的长度 if err = writeUint64(e.bw, lenField, e.uint64buf); err != nil { return err } // 涉及对齐时追加padding if padBytes != 0 { data = append(data, make([]byte, padBytes)...) } // 写入recode内容 n, err = e.bw.Write(data) // 记录写入量 walWriteBytes.Add(float64(n)) return err }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37Record序列化后会以Frame的格式持久化。Frame头部是一个64bit长度字段,其中MSB标识这个Frame是否有padding字节,接下来才是真正序列化的数据长度。func encodeFrameSize(dataBytes int) (lenField uint64, padBytes int) { // 数据长度 lenField = uint64(dataBytes) // 补齐的padding长度 padBytes = (8 - (dataBytes % 8)) % 8 // 高56位记录padding长度 if padBytes != 0 { // 最高位是1标识含有padding lenField |= uint64(0x80|padBytes) << 56 } return lenField, padBytes }1
2
3
4
5
6
7
8
9
10
11
12
# 2.2.WAL创建
wal.Create()会创建WAL实例、写入初始元数据、初始化编解码器等,用于后续对WAL文件的追加或读取操作。func Create(lg *zap.Logger, dirpath string, metadata []byte) (*WAL, error) { // 先基于临时目录初始化 tmpdirpath := filepath.Clean(dirpath) + ".tmp" ... // 创建文件互斥锁 p := filepath.Join(tmpdirpath, walName(0, 0)) f, err := fileutil.LockFile(p, os.O_WRONLY|os.O_CREATE, fileutil.PrivateFileMode) ... // 定位到文件末尾 if _, err = f.Seek(0, io.SeekEnd); err != nil { return nil, err } // 预分配文件,大小为64MB if err = fileutil.Preallocate(f.File, SegmentSizeBytes, true); err != nil { return nil, err } // 初始化WAL实例 w := &WAL{ lg: lg, dir: dirpath, metadata: metadata, } // 初始化编码器 w.encoder, err = newFileEncoder(f.File, 0) // 上互斥锁的文件加入locks数组 w.locks = append(w.locks, f) // 存储校验码 if err = w.saveCrc(0); err != nil { return nil, err } // 将元数据的Record记录到WAL的header处 if err = w.encoder.encode(&walpb.Record{Type: metadataType, Data: metadata}); err != nil { return nil, err } // 保存空的快照信息 if err = w.SaveSnapshot(walpb.Snapshot{}); err != nil { return nil, err } // 临时文件重命名 if w, err = w.renameWAL(tmpdirpath); err != nil { return nil, err } var perr error defer func() { if perr != nil { w.cleanupWAL(lg) } }() // 打开父目录 pdir, perr := fileutil.OpenDir(filepath.Dir(w.dir)) if perr != nil { return nil, perr } // 刷盘,确保目录重命名就及元数据生效 if perr = fileutil.Fsync(pdir); perr != nil { return nil, perr } // 关闭父目录描述符 if perr = pdir.Close(); perr != nil { return nil, perr } return w, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
651.预分配
preallocate预先分配一块固定大小的文件,通过fcntl的对应flag实现,实际上这个预先分配的动作很耗时,直接append效果应该是差不多的,但硬盘不够的时候可能最后会出现没法写日志的情况,所以会利用filePipeline利用协程去预先申请空间。2.Fsync
Linux系统上write不会立即写入磁盘,直到缓冲区满后才会同步到磁盘,增加效率 。Fsync可以让缓存的数据立刻刷盘,避免etcd共识机制要求下日志丢失,raft状态恢复失败。当然,除了fsync还有一个fdatasync,可以避免更新不必要的信息,比如文件时间之类的。
# 2.3.WAL存储
WAL主要用来持久化存储日志,raft收到一个proposal时会调用Save方法完成持久化。func (w *WAL) Save(st raftpb.HardState, ents []raftpb.Entry) error { w.mu.Lock() defer w.mu.Unlock() // raft硬状态为空同时没有需要持久化的存储日志 if raft.IsEmptyHardState(st) && len(ents) == 0 { return nil } // 是否需要刷盘 mustSync := raft.MustSync(st, w.state, len(ents)) // 保存日志项 for i := range ents { if err := w.saveEntry(&ents[i]); err != nil { return err } } // 持久化hardState if err := w.saveState(&st); err != nil { return err } // 获取最后一个LockedFile的大小(正在使用) curOff, err := w.tail().Seek(0, io.SeekCurrent) ... // 最后一个LockedFile文件小于64MB if curOff < SegmentSizeBytes { // 刷盘 if mustSync { // 写刷盘数据 err = w.sync() // gofail: var walAfterSync struct{} return err } return nil } // WAL文件超出64MB,切割 return w.cut() }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40MustSync用来判断当前的Save是否需要同步持久化,由于每台服务器都必须无条件持久化三个量—currentTerm、votedFor和logEntries,因此只要任意一个改变就需要持久化。func MustSync(st, prevst pb.HardState, entsnum int) bool { return entsnum != 0 || st.Vote != prevst.Vote || st.Term != prevst.Term }1
2
3日志
Entry的持久化由saveEntry完成,会先封装成一个Record,其次encode持久化。func (w *WAL) saveEntry(e *raftpb.Entry) error { // 日志数据序列化为json b := pbutil.MustMarshal(e) // 初始化日志类型record rec := &walpb.Record{Type: entryType, Data: b} // 编码写入buf if err := w.encoder.encode(rec); err != nil { return err } // 记录当前持久化索引 w.enti = e.Index return nil }1
2
3
4
5
6
7
8
9
10
11
12
13hardState的持久化由saveState完成,依然封装一个Record然后编码写入buf缓冲。func (w *WAL) saveState(s *raftpb.HardState) error { // 记录state w.state = *s // 序列化为json b := pbutil.MustMarshal(s) // 初始化为record rec := &walpb.Record{Type: stateType, Data: b} // 编码写入buf缓冲 return w.encoder.encode(rec) }1
2
3
4
5
6
7
8
9
10WAL超出一定大小时会进行切割,默认限制大小为64MB,主流程由cut实现。func (w *WAL) cut() error { // 获取当前写入文件的偏移量 off, serr := w.tail().Seek(0, io.SeekCurrent) ... // 截断文件 if err := w.tail().Truncate(off); err != nil { return err } // 刷盘,已写入文件持久化 if err := w.sync(); err != nil { return err } // 准备新文件 fpath := filepath.Join(w.dir, walName(w.seq()+1, w.enti+1)) // 基于预分配机制打开一个临时文件 newTail, err := w.fp.Open() ... // 新文件添加到LockedFile数组 w.locks = append(w.locks, newTail) // 获取之前的crc prevCrc := w.encoder.crc.Sum32() // 初始化新文件的encoder,继承之前使用的crc,确保跨文件校验连续性 w.encoder, err = newFileEncoder(w.tail().File, prevCrc) ... // 保存crc类型的record if err = w.saveCrc(prevCrc); err != nil { return err } // 保存metadata类型的record if err = w.encoder.encode(&walpb.Record{Type: metadataType, Data: w.metadata}); err != nil { return err } // 保存hardState类型的record if err = w.saveState(&w.state); err != nil { return err } // 刷盘 if err = w.sync(); err != nil { return err } // 找到当前正在使用文件的写入offset off, err = w.tail().Seek(0, io.SeekCurrent) if err != nil { return err } // 临时文件重命名 if err = os.Rename(newTail.Name(), fpath); err != nil { return err } // 刷盘,确保操作持久化 if err = fileutil.Fsync(w.dirFile); err != nil { return err } ... // 关闭临时文件 newTail.Close() // 重新打开 if newTail, err = fileutil.LockFile(fpath, os.O_WRONLY, fileutil.PrivateFileMode); err != nil { return err } // 根据偏移量指向下一块写入位置 if _, err = newTail.Seek(off, io.SeekStart); err != nil { return err } // 更新LockedFile数组的文件句柄 w.locks[len(w.locks)-1] = newTail // 更新重命名文件的encoder prevCrc = w.encoder.crc.Sum32() w.encoder, err = newFileEncoder(w.tail().File, prevCrc) if err != nil { return err } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77流程
w.tail().Truncate() → w.sync() → filePipeline.Open() → 写入元数据 → w.sync() → os.Rename() → w.fsync() → fileutil.LockFile()
# 2.4.WAL开启
Open()会根据对应的index读取后面的所有日志。func Open(lg *zap.Logger, dirpath string, snap walpb.Snapshot) (*WAL, error) { // 只打开最后一个seq小于snap中的index之后的所有WAL文件(写模式) w, err := openAtIndex(lg, dirpath, snap, true) ... // 打开WAL文件 if w.dirFile, err = fileutil.OpenDir(w.dir); err != nil { return nil, err } return w, nil } func OpenForRead(lg *zap.Logger, dirpath string, snap walpb.Snapshot) (*WAL, error) { // 只读打开 return openAtIndex(lg, dirpath, snap, false) } // 打开指定位置的WAL文件 func openAtIndex(lg *zap.Logger, dirpath string, snap walpb.Snapshot, write bool) (*WAL, error) { // 获取所有WAL名称,需要读取的最小nameIndex names, nameIndex, err := selectWALFiles(lg, dirpath, snap) ... // 根据指定位置的WAL文件索引,打开后面所有文件 rs, ls, closer, err := openWALFiles(lg, dirpath, names, nameIndex, write) ... // 初始化一个WAL对象 w := &WAL{ lg: lg, dir: dirpath, start: snap, decoder: newDecoder(rs...), readClose: closer, locks: ls, } // 写模式,执行预分配 if write { ... if _, _, err := parseWALName(filepath.Base(w.tail().Name())); err != nil { closer() return nil, err } // 一直执行预分配等待消费 w.fp = newFilePipeline(lg, w.dir, SegmentSizeBytes) } return w, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
# 2.5.WAL读取
ReadAll()用于读取WAL文件内容,其实就是decode解码出Open()的文件内容,解析到结构体中使用。func (w *WAL) ReadAll() (metadata []byte, state raftpb.HardState, ents []raftpb.Entry, err error) { w.mu.Lock() defer w.mu.Unlock() rec := &walpb.Record{} ... decoder := w.decoder var match bool // decode循环解码record,直到遇到错误 for err = decoder.decode(rec); err == nil; err = decoder.decode(rec) { switch rec.Type { // 日志类型 case entryType: // 反序列化日志内容 e := mustUnmarshalEntry(rec.Data) // 0 <= e.Index-w.start.Index - 1 < len(ents) // 日志条目的索引大于WAL应该读取的其实Index if e.Index > w.start.Index { // 根据偏移量计算具体索引 up := e.Index - w.start.Index - 1 ... // 未越界(连续)追加 ents = append(ents[:up], e) } w.enti = e.Index // hardState类型 case stateType: // 反序列化内容 state = mustUnmarshalState(rec.Data) // 元数据类型 case metadataType: // 一致性检查,确保所有元数据记录相同 if metadata != nil && !bytes.Equal(metadata, rec.Data) { // 不相同时重置 state.Reset() return nil, state, nil, ErrMetadataConflict } metadata = rec.Data // crc类型 case crcType: // 计算当前crc crc := decoder.crc.Sum32() // 一致性检查 ... // 更新crc decoder.updateCRC(rec.Crc) // 快照类型 case snapshotType: var snap walpb.Snapshot // 反序列化数据 pbutil.MustUnmarshal(&snap, rec.Data) ... // 一致性检查 ... } } // 查找最后使用的lockedFile switch w.tail() { ... default: ... // 获取当前活跃的WAL文件(最后一个段) // lastOffset()指向解码器最后成功解码的位置 // io.SeekStart代表文件开头计算偏移量 // 文件指针移动到最后一个有效记录末尾 if _, err = w.tail().Seek(w.decoder.lastOffset(), io.SeekStart); err != nil { return nil, state, nil, err } // 清零最后一次解析成功位置后的数据,避免部分写入的脏数据 if err = fileutil.ZeroToEnd(w.tail().File); err != nil { return nil, state, nil, err } } ... // 记录元数据 w.metadata = metadata // 初始化encoder if w.tail() != nil { // create encoder (chain crc with the decoder), enable appending w.encoder, err = newFileEncoder(w.tail().File, w.decoder.lastCRC()) ... } w.decoder = nil return metadata, state, ents, err }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
# 2.6.FilePipeline
文件
Pipeline采用饿汉式思想,提前创建备用文件,加快文件的创建速度。type filePipeline struct { dir string // 存放临时文件的目录 size int64 // 创建临时文件时预分配空间的大小,默认64MB count int // 当前filePipeline实例创建的临时文件数 filec chan *fileutil.LockedFile // 新建的临时文件句柄会通过filec通道返回给WAL实例使用 errc chan error // 创建临时文件出现异常时放入errc通道 donec chan struct{} // filePipeline.Close()调用时关闭donec通知filePipeline删除最周创建的临时文件 }1
2
3
4
5
6
7
8
9
10
11
12
13newFilePipeline()创建filePipeline实例时会启动一个后台协程执行filePipeline.run()方法,该方法会创建新的临时文件将其句柄传递到filec通道。func newFilePipeline(lg *zap.Logger, dir string, fileSize int64) *filePipeline { fp := &filePipeline{ lg: lg, dir: dir, size: fileSize, filec: make(chan *fileutil.LockedFile), errc: make(chan error, 1), donec: make(chan struct{}), } // run()方法内申请临时文件并缓存至管道 go fp.run() return fp } func (fp *filePipeline) run() { defer close(fp.errc) for { // 申请文件 f, err := fp.alloc() if err != nil { fp.errc <- err return } select { // 存入filec管道待使用 case fp.filec <- f: // filePipeline关闭 case <-fp.donec: // 移除最后一次创建文件 os.Remove(f.Name()) f.Close() return } } } // 申请临时文件 func (fp *filePipeline) alloc() (f *fileutil.LockedFile, err error) { // 创建临时文件的编号0或1 fpath := filepath.Join(fp.dir, fmt.Sprintf("%d.tmp", fp.count%2)) // 指定文件创建时的模式和权限 if f, err = fileutil.LockFile(fpath, os.O_CREATE|os.O_WRONLY, fileutil.PrivateFileMode); err != nil { return nil, err } // 尝试预分配,当前文件系统不支持时不报错 if err = fileutil.Preallocate(f.File, fp.size, true); err != nil { // 预分配错误关闭文件 f.Close() return nil, err } // 预分配成功更新临时文件数 fp.count++ return f, nil } // 关闭filePipeline实例的donec管道 func (fp *filePipeline) Close() error { close(fp.donec) return <-fp.errc }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
