分布式共识算法Raft依赖日志复制来保证多个节点状态一致,而日志本身必须持久化到磁盘,否则节点重启后会丢失已提交的记录,破坏安全性。在C#中实现Raft持久化日志,核心要解决三个问题:日志如何编码落盘、如何高效追加与读取、如何在崩溃后安全恢复。本文围绕这几个角度展开,给出可用的代码结构与设计权衡。

日志文件结构与编码设计
Raft日志由一条条有序记录组成,每条记录至少包含索引(Index)、任期(Term)、命令类型与具体数据。为了能在C#中稳定读写,我们通常采用二进制格式而非文本格式,因为二进制可以避免字符编码带来的歧义,也更容易计算校验和。一个常见的设计是:文件头部写入魔数(Magic Number)和版本号,随后每一笔日志以「长度前缀+内容+CRC32」的形式追加。
长度前缀使用固定4字节整数,内容部分依次写入Index(8字节)、Term(8字节)、CommandType(1字节)、Payload长度(4字节)与Payload字节数组。末尾的CRC32覆盖前面所有字节,读取时重新计算并比对,若不一致说明写中断或磁盘损坏。这种结构在C#里用BinaryWriter实现非常直接,而且由于是纯追加写,不需要修改已落盘区域,降低了并发复杂度。
值得注意的是,Raft要求日志顺序严格递增,因此在编码时必须保证Index连续。如果某次写入因进程崩溃只写了半条记录,恢复阶段应通过CRC失败识别并截断该半条,而不是尝试解析残缺数据。下面的代码展示了单条日志的编码方法:
using System;
using System.IO;
using System.IO.Compression;
using System.Security.Cryptography;
public class RaftLogEntry
{
public long Index;
public long Term;
public byte CommandType;
public byte[] Payload;
}
public static class RaftLogEncoder
{
public static void WriteEntry(BinaryWriter writer, RaftLogEntry entry)
{
using (var ms = new MemoryStream())
{
var w = new BinaryWriter(ms);
w.Write(entry.Index);
w.Write(entry.Term);
w.Write(entry.CommandType);
w.Write(entry.Payload.Length);
w.Write(entry.Payload);
w.Flush();
byte[] body = ms.ToArray();
uint crc = Crc32.Compute(body);
writer.Write(body.Length);
writer.Write(body);
writer.Write(crc);
writer.Flush();
}
}
}
public static class Crc32
{
private static uint[] table;
static Crc32()
{
table = new uint[256];
for (uint i = 0; i < 256; i++)
{
uint c = i;
for (int k = 0; k < 8; k++)
{
c = ((c & 1) != 0) ? (0xEDB88320 ^ (c >> 1)) : (c >> 1);
}
table[i] = c;
}
}
public static uint Compute(byte[] data)
{
uint crc = 0xFFFFFFFF;
foreach (byte b in data)
{
crc = table[(crc ^ b) & 0xFF] ^ (crc >> 8);
}
return crc ^ 0xFFFFFFFF;
}
}
追加写入与刷盘策略
在C#里实现追加写,最简单的方式是以FileMode.Append打开文件,并用BufferedStream包一层BinaryWriter。但Raft算法对持久性有明确要求:领导者在回应客户端之前,必须确保日志已落到多数节点磁盘。本地节点至少也要调用Flush让操作系统把数据交到磁盘控制器,否则断电仍可能丢数据。因此写完后应调用FileStream.Flush(true),而不是仅清空托管缓冲区。
另一个常见问题是写入性能。如果每条日志都同步刷盘,吞吐会受限于磁盘IOPS。实际系统中可引入组提交(Group Commit),即在内存中积攒几毫秒内的日志,合并一次刷盘。C#的Channel或BlockingCollection能方便实现生产者消费者模型,写线程把日志放入队列,落盘线程批量取出并写入。这样既保留持久化语义,又显著提升吞吐。下面的示例展示了一个简单的异步批量写入器骨架:
using System;
using System.Collections.Concurrent;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
public class BatchRaftLogWriter : IDisposable
{
private readonly BlockingCollection<RaftLogEntry> queue =
new BlockingCollection<RaftLogEntry>(1024);
private readonly FileStream fs;
private readonly BinaryWriter writer;
private readonly Thread worker;
private volatile bool running = true;
public BatchRaftLogWriter(string path)
{
fs = new FileStream(path, FileMode.Append, FileAccess.Write,
FileShare.Read, 4096, FileOptions.WriteThrough);
writer = new BinaryWriter(fs);
worker = new Thread(Loop) { IsBackground = true };
worker.Start();
}
public void Append(RaftLogEntry entry)
{
queue.Add(entry);
}
private void Loop()
{
while (running || queue.Count > 0)
{
var batch = new System.Collections.Generic.List<RaftLogEntry>();
while (batch.Count < 64 && queue.TryTake(out var e))
{
batch.Add(e);
}
if (batch.Count == 0)
{
Thread.Sleep(1);
continue;
}
foreach (var e in batch)
{
RaftLogEncoder.WriteEntry(writer, e);
}
fs.Flush(true);
}
}
public void Dispose()
{
running = false;
queue.CompleteAdding();
worker.Join();
writer.Dispose();
fs.Dispose();
}
}
对于极高吞吐场景,也可以考虑使用MemoryMappedFile,把文件映射到进程地址空间,由操作系统负责页回写。但这种方式在崩溃恢复时边界控制更复杂,且C#的封装不如BinaryWriter直观,中小规模集群用 buffered 追加写已经足够。选择方案时应权衡团队对底层API的熟悉度与运维排错成本。
崩溃恢复与日志截断
节点重启后第一件事是重放日志文件。恢复流程必须按物理顺序逐条读取,利用前面写的CRC校验每条完整性。一旦碰到长度前缀声明的大小和实际剩余字节不符,或CRC不匹配,就说明文件在这之后损坏,应当停止读取并将文件截断到上一条正确记录的末尾。C#里可以用FileStream的SetLength方法完成截断,但要注意先刷新再截断,避免操作系统缓存导致截断位置错误。
除了崩溃恢复,Raft还允许通过快照压缩旧日志。当某节点状态机已应用前N条日志,就可以把N之前的记录从文件删除,只保留一个「快照点索引」作为新起点。实现上可另建一个meta文件记录快照点,加载时先读meta,再从对应Index之后的日志继续。这样日志文件不会无限膨胀,也缩短了恢复时间。下面代码演示了顺序加载并校验的过程:
using System;
using System.IO;
using System.Collections.Generic;
public static class RaftLogRecovery
{
public static List<RaftLogEntry> Load(string path)
{
var result = new List<RaftLogEntry>();
using (var fs = new FileStream(path, FileMode.OpenOrCreate,
FileAccess.Read, FileShare.Read))
using (var reader = new BinaryReader(fs))
{
while (fs.Position < fs.Length)
{
long entryStart = fs.Position;
try
{
int len = reader.ReadInt32();
byte[] body = reader.ReadBytes(len);
if (body.Length < len)
break;
uint crc = reader.ReadUInt32();
if (Crc32.Compute(body) != crc)
break;
using (var ms = new MemoryStream(body))
using (var r = new BinaryReader(ms))
{
var e = new RaftLogEntry();
e.Index = r.ReadInt64();
e.Term = r.ReadInt64();
e.CommandType = r.ReadByte();
int plen = r.ReadInt32();
e.Payload = r.ReadBytes(plen);
result.Add(e);
}
}
catch (Exception)
{
fs.Position = entryStart;
fs.SetLength(entryStart);
break;
}
}
}
return result;
}
}
恢复逻辑里还有一个细节:必须检查加载出的日志Index是否连续,且Term是否符合单调或重置规则。若发现某条Index与前一条不连续,说明写入时发生了严重错误,应报错而不是默默采用。只有经过严格校验的日志,才能交给Raft状态机去匹配投票与提交逻辑,从而保障分布式系统的整体正确性。