RocketMQ 刷盘机制:同步、异步与顺序写的秘密
一句话结论(30s)
RocketMQ 刷盘机制本质是「可靠性 vs 吞吐量」的取舍:同步刷盘必须等 fsync() 真正落盘才返回 ACK(消息不丢但吞吐低),异步刷盘写 Page Cache 即 ACK(吞吐高但 OS 崩溃会丢)。因为底层用 mmap + 顺序写把随机写提速 100 倍,再用 GroupCommit 批量 fsync 摊薄同步刷盘开销,才能在普通磁盘跑到数十万 TPS。
核心原理(2min)
消息经 mmap 直接写入 CommitLog 映射区域(Page Cache),刷盘策略决定 ACK 时机——SYNC_FLUSH 等 fsync() 落盘后 ACK,ASYNC_FLUSH(默认)写 Page Cache 后立即 ACK,后台线程按 flushIntervalCommitLog(默认 500ms)批量刷盘。关键机制:GroupCommit 把并发到达的一批请求聚合成一次 fsync(),大幅提升同步刷盘吞吐;ConsumeQueue 是冗余索引,刷盘间隔(1s)比 CommitLog 更宽松,丢了可从 CommitLog 重建。
底层深入(5-10min)
一个关键问题:消息什么时候算”持久化成功”?
RocketMQ 收到一条消息,写入 CommitLog 后,不一定立即持久化到磁盘。“写入成功”的定义取决于刷盘策略。 这个定义直接影响消息的可靠性和吞吐量。
两种刷盘模式
同步刷盘(SYNC_FLUSH)
同步刷盘的流程:
生产者发送消息
│
▼
Broker 接收 → 写入内存(Page Cache)
│
▼
调用 fsync() → 强制刷盘 → 数据写入磁盘
│
▼
返回 ACK 给生产者
消息确认时机:必须在数据真正落盘之后。如果 fsync() 还没完成就宕机,消息丢失。
优缺点:
| 优点 | 缺点 |
|---|---|
| 消息可靠性最高,数据不丢失 | 吞吐量低(受限于磁盘顺序写速度,通常几百到几千 TPS) |
| 金融、支付等强一致性场景适用 | 每次写入都要等磁盘 IO 完成,延迟高 |
关键在于:fsync() 是一次系统调用,它告诉操作系统”把 Page Cache 中该文件的数据立即写入磁盘”。同步刷盘时,RocketMQ 在返回 ACK 之前必须等待 fsync() 完成。
异步刷盘(ASYNC_FLUSH,默认)
异步刷盘的流程:
生产者发送消息
│
▼
Broker 接收 → 写入内存(Page Cache / mmap 映射区域)→ 立即返回 ACK
│
▼
后台线程定时调用 fsync() → 批量刷盘
消息确认时机:数据写入 Page Cache 后立即返回 ACK。这时数据还没有真正落盘——它在操作系统的内存缓冲区中。如果操作系统崩溃,Page Cache 中的数据会丢失。
优缺点:
| 优点 | 缺点 |
|---|---|
| 吞吐量高(可达数万到十万 TPS) | 极端情况下可能丢消息(操作系统崩溃) |
| 延迟低(不需要等磁盘 IO) | Broker 进程崩溃一般不丢(OS 会继续刷盘) |
💭 思考:同样是异步刷盘,为什么「Broker 进程崩溃」和「操作系统崩溃」的后果不一样?—— 因为数据写进 Page Cache 后,归属权已经交给了 OS 内核:进程挂了,OS 还在、会继续把脏页刷到磁盘,所以 Broker 崩溃一般不丢;但若整个 OS 崩溃(掉电/内核 panic),Page Cache 里还没落盘的部分就永久没了。所以异步刷盘的「丢」是针对 OS 级故障,而不是进程级故障——这决定了它适合「允许极少丢失」的场景。
关键配置
# broker.properties
flushDiskType = ASYNC_FLUSH # 异步刷盘(默认)
# flushDiskType = SYNC_FLUSH # 同步刷盘
# 异步刷盘相关参数(仅在 ASYNC_FLUSH 下生效)
flushIntervalCommitLog = 500 # 每隔 500ms 执行一次 CommitLog 刷盘
commitIntervalCommitLog = 200 # CommitLog 页提交间隔(ms)
flushCommitLogTimed = false # true = 定时,false = 实时(有新消息就触发检查)
# 同步刷盘参数
syncFlushTimeout = 5000 # 同步刷盘超时时间(ms)
内存映射文件(mmap):零拷贝的存储基础
RocketMQ 使用 内存映射文件(Memory-Mapped File) 来操作 CommitLog。这是理解刷盘机制的前提。
传统写文件 vs mmap
传统 write():
用户态 Buffer → write() 系统调用 → 内核拷贝到 Page Cache → 刷盘
(两次拷贝:用户态 → 内核,内核 → 磁盘)
mmap:
文件映射到虚拟地址空间 → 直接写入映射区域 → Page Cache → 刷盘
(一次拷贝:用户态直接操作 Page Cache,绕过了用户态 Buffer)
真实源码:DefaultMappedFile 的 mmap 与顺序追加写
RocketMQ 里真正干活的类是 DefaultMappedFile(MappedFile 只是接口)。核心就两步:init 时 mmap 映射文件,appendMessagesInner 时把消息顺序写进映射内存。
// org.apache.rocketmq.store.logfile.DefaultMappedFile#init —— 用 mmap 把文件映射进进程地址空间
private void init(final String fileName, final int fileSize, final RunningFlags runningFlags) throws IOException {
this.fileName = fileName;
this.fileSize = fileSize;
this.file = new File(fileName);
this.fileFromOffset = Long.parseLong(this.file.getName());
...
try {
this.fileChannel = new RandomAccessFile(this.file, "rw").getChannel();
if (writeWithoutMmap) {
// Still create MappedByteBuffer for reading operations
this.mappedByteBuffer = this.fileChannel.map(MapMode.READ_ONLY, 0, fileSize);
} else {
// Use MappedByteBuffer for both reading and writing (default behavior)
this.mappedByteBuffer = this.fileChannel.map(MapMode.READ_WRITE, 0, fileSize);
}
TOTAL_MAPPED_VIRTUAL_MEMORY.addAndGet(fileSize);
...
} catch (IOException e) {
log.error("Failed to map file " + this.fileName, e);
throw e;
} finally {
...
}
}
fileChannel.map(MapMode.READ_WRITE, 0, fileSize) 这一步就是把 1GB 文件映射成 MappedByteBuffer——之后写入这个 buffer 就等于直接操作 Page Cache,绕过了 write() 系统调用和用户态 Buffer 的一次拷贝。
// DefaultMappedFile#appendMessagesInner —— 顺序追加写:切片定位到 wrotePosition,写完推进写指针
public AppendMessageResult appendMessagesInner(final MessageExt messageExt, final AppendMessageCallback cb,
PutMessageContext putMessageContext) {
...
int currentPos = WROTE_POSITION_UPDATER.get(this);
long fileFromOffset = this.getFileFromOffset();
if (currentPos < this.fileSize) {
...
ByteBuffer byteBuffer = appendMessageBuffer().slice();
byteBuffer.position(currentPos);
AppendMessageResult result;
try {
...
result = cb.doAppend(fileFromOffset, byteBuffer, this.fileSize - currentPos,
(MessageExtBrokerInner) messageExt, putMessageContext);
...
} finally {
...
}
WROTE_POSITION_UPDATER.addAndGet(this, result.getWroteBytes());
...
return result;
}
log.error("MappedFile.appendMessage return null, wrotePosition: {} fileSize: {}", currentPos, this.fileSize);
return new AppendMessageResult(AppendMessageStatus.UNKNOWN_ERROR);
}
分析:写消息前先取 wrotePosition(写指针),把 mappedByteBuffer.slice() 定位到该 offset 后交给 doAppend 编码写入,写完用 WROTE_POSITION_UPDATER.addAndGet 推进写指针。写指针永远单调递增、只往末尾追加,这就是「顺序写」在代码层面的实现——没有 seek 到中间覆盖任何旧数据。
顺序写 = 快
为什么 RocketMQ 的写性能能达到数十万 TPS?核心原因就是顺序写磁盘:
RocketMQ 顺序写(Append-only):
磁盘磁头不动,数据连续写入相邻扇区
速度:600MB/s(SSD)~ 100MB/s(HDD)
随机写(MySQL InnoDB 的 B+Tree 更新):
磁头不断移动,写入分散的扇区
速度:几 MB/s(HDD)
结论:顺序写比随机写快 100 倍以上
RocketMQ 把所有 Topic 的消息全部追加到一个 CommitLog 文件末尾,把随机写转化成了顺序写。这是 Kafka 和 RocketMQ 高吞吐的基石。
🤔 思考穿插:为什么「顺序写」能比随机写快 100 倍?—— 因为磁盘上最贵的操作是磁头寻道,顺序追加让磁头基本不动、数据连续写入相邻扇区;随机写则要不断跳转寻道。所以 RocketMQ/Kafka 都把散落的随机写收敛成「Append-only 顺序写」,这是它们高吞吐的第一性原理。
刷盘线程:异步刷盘的调度细节
CommitLog 刷盘决策:handleDiskFlush(同步/异步)
消息写完 CommitLog 后,CommitLog 调 flushManager.handleDiskFlush 决定「什么时候 ACK」。同步还是异步,只看一个枚举:
// org.apache.rocketmq.store.config.FlushDiskType —— 就两个值
public enum FlushDiskType {
SYNC_FLUSH,
ASYNC_FLUSH
}
// CommitLog.FlushManager#handleDiskFlush —— 同步走 GroupCommitService,异步只唤醒后台线程
@Override
public CompletableFuture<PutMessageStatus> handleDiskFlush(AppendMessageResult result, MessageExt messageExt) {
// Synchronization flush
if (FlushDiskType.SYNC_FLUSH == CommitLog.this.defaultMessageStore.getMessageStoreConfig().getFlushDiskType()) {
final GroupCommitService service = (GroupCommitService) this.flushCommitLogService;
if (messageExt.isWaitStoreMsgOK()) {
GroupCommitRequest request = new GroupCommitRequest(result.getWroteOffset() + result.getWroteBytes(), CommitLog.this.defaultMessageStore.getMessageStoreConfig().getSyncFlushTimeout());
flushDiskWatcher.add(request);
service.putRequest(request);
return request.future();
} else {
service.wakeup();
return CompletableFuture.completedFuture(PutMessageStatus.PUT_OK);
}
}
// Asynchronous flush
else {
if (!CommitLog.this.defaultMessageStore.isTransientStorePoolEnable()) {
if (defaultMessageStore.getMessageStoreConfig().isWakeFlushWhenPutMessage()) {
flushCommitLogService.wakeup();
}
} else {
if (defaultMessageStore.getMessageStoreConfig().isWakeCommitWhenPutMessage()) {
commitRealTimeService.wakeup();
}
}
return CompletableFuture.completedFuture(PutMessageStatus.PUT_OK);
}
}
分析:同步刷盘把请求包装成 GroupCommitRequest 丢给 GroupCommitService 并返回 request.future()——生产者的 ACK 会一直阻塞在这个 future 上,直到刷盘线程把数据刷过该 offset 才被唤醒。异步刷盘则直接 return PUT_OK,只 wakeup() 后台刷盘线程,ACK 完全不等待磁盘 IO。
💭 思考:为什么刷盘策略要收敛成
SYNC_FLUSH/ASYNC_FLUSH一个枚举,让handleDiskFlush统一决定 ACK 时机,而不是在写消息处各写一套逻辑?—— 因为「什么时候 ACK」是刷盘机制唯一的决策点:同步刷盘把 ACK 挂到future上等落盘,异步刷盘立刻PUT_OK。收敛到一个枚举 + 一个决策方法,写路径就能共用同一套代码,切换刷盘策略只改配置不改逻辑,也把「可靠性 vs 吞吐」的取舍显式成了一个开关。
底层落盘:flush 与 commit
handleDiskFlush 只是「决策」,真正把数据刷到磁盘的是 DefaultMappedFile#flush——mappedByteBuffer.force() 底层就是 fsync():
// DefaultMappedFile#flush —— mappedByteBuffer.force() 即 fsync
public int flush(final int flushLeastPages) {
if (!isWriteable()) {
return this.getFlushedPosition();
}
if (this.isAbleToFlush(flushLeastPages)) {
if (this.hold()) {
int value = getReadPosition();
try {
this.mappedByteBufferAccessCountSinceLastSwap++;
//We only append data to fileChannel or mappedByteBuffer, never both.
if (writeWithoutMmap || writeBuffer != null || this.fileChannel.position() != 0) {
this.fileChannel.force(false);
} else {
this.mappedByteBuffer.force();
}
this.lastFlushTime = System.currentTimeMillis();
FLUSHED_POSITION_UPDATER.set(this, value);
} catch (Throwable e) {
if (e instanceof IOException) {
getAndMakeNotWriteable();
}
log.error("Error occurred when force data to disk.", e);
}
this.release();
} else {
log.warn("in flush, hold failed, flush offset = " + FLUSHED_POSITION_UPDATER.get(this));
FLUSHED_POSITION_UPDATER.set(this, getReadPosition());
}
}
return this.getFlushedPosition();
}
分析:mappedByteBuffer.force() 是 Java 封装的 fsync,把 Page Cache 里已 wrote 但未 flushed 的数据真正刷到磁盘,然后把 FLUSHED_POSITION 推进到刷盘位置。isAbleToFlush(flushLeastPages) 就是「脏页攒够 N 页才刷」的闸门,避免每条消息都触发一次昂贵的 fsync。
还有一个 commit 只在你开启 transientStorePoolEnable(堆外内存池)时才生效:写入先落到堆外 writeBuffer,再由 commit 线程批量提交到 FileChannel 的 Page Cache:
// DefaultMappedFile#commit / #commit0 —— 把堆外 writeBuffer 提交到 FileChannel
public int commit(final int commitLeastPages) {
if (writeBuffer == null) {
//no need to commit data to file channel, so just regard wrotePosition as committedPosition.
return WROTE_POSITION_UPDATER.get(this);
}
...
if (transientStorePool != null && !transientStorePool.isRealCommit()) {
COMMITTED_POSITION_UPDATER.set(this, WROTE_POSITION_UPDATER.get(this));
} else if (this.isAbleToCommit(commitLeastPages)) {
if (this.hold()) {
commit0();
this.release();
} else {
log.warn("in commit, hold failed, commit offset = " + COMMITTED_POSITION_UPDATER.get(this));
}
}
...
return COMMITTED_POSITION_UPDATER.get(this);
}
protected void commit0() {
int writePos = WROTE_POSITION_UPDATER.get(this);
int lastCommittedPosition = COMMITTED_POSITION_UPDATER.get(this);
if (writePos - lastCommittedPosition > 0) {
try {
ByteBuffer byteBuffer = writeBuffer.slice();
byteBuffer.position(lastCommittedPosition);
byteBuffer.limit(writePos);
this.fileChannel.position(lastCommittedPosition);
this.fileChannel.write(byteBuffer);
COMMITTED_POSITION_UPDATER.set(this, writePos);
} catch (Throwable e) {
log.error("Error occurred when commit data to FileChannel.", e);
}
}
}
分析:commit0 把堆外 writeBuffer 从 lastCommittedPosition 到 writePos 这一段 fileChannel.write 进 Page Cache,再推进 COMMITTED_POSITION。wrotePosition → committedPosition → flushedPosition 三级游标,正是 RocketMQ 区分「已写入 / 已提交到 PageCache / 已落盘」三个阶段的依据。
🤔 思考穿插:为什么 RocketMQ 要区分 wrote / committed / flushed 三级游标,而不是一个 offset 走天下?—— 因为「写入进程内存」「提交到 PageCache」「真正落盘」是物理上不同的三个阶段,异步刷盘下它们之间存在时间差;只有分开记录,才能在故障恢复时判断数据到底走到了哪一步、会不会丢。
CommitLog 异步刷盘(FlushRealTimeService)
// CommitLog.FlushRealTimeService#run —— 异步刷盘后台线程
@Override
public void run() {
CommitLog.log.info("{} service started", this.getServiceName());
while (!this.isStopped()) {
boolean flushCommitLogTimed = CommitLog.this.defaultMessageStore.getMessageStoreConfig().isFlushCommitLogTimed();
int interval = CommitLog.this.defaultMessageStore.getMessageStoreConfig().getFlushIntervalCommitLog();
int flushPhysicQueueLeastPages = CommitLog.this.defaultMessageStore.getMessageStoreConfig().getFlushCommitLogLeastPages();
int flushPhysicQueueThoroughInterval =
CommitLog.this.defaultMessageStore.getMessageStoreConfig().getFlushCommitLogThoroughInterval();
boolean printFlushProgress = false;
// Print flush progress
long currentTimeMillis = System.currentTimeMillis();
if (currentTimeMillis >= (this.lastFlushTimestamp + flushPhysicQueueThoroughInterval)) {
this.lastFlushTimestamp = currentTimeMillis;
flushPhysicQueueLeastPages = 0;
printFlushProgress = (printTimes++ % 10) == 0;
}
try {
if (flushCommitLogTimed) {
Thread.sleep(interval);
} else {
this.waitForRunning(interval);
}
if (printFlushProgress) {
this.printFlushProgress();
}
long begin = System.currentTimeMillis();
CommitLog.this.mappedFileQueue.flush(flushPhysicQueueLeastPages);
long storeTimestamp = CommitLog.this.mappedFileQueue.getStoreTimestamp();
if (storeTimestamp > 0) {
CommitLog.this.defaultMessageStore.getStoreCheckpoint().setPhysicMsgTimestamp(storeTimestamp);
}
long past = System.currentTimeMillis() - begin;
CommitLog.this.getMessageStore().getPerfCounter().flowOnce("FLUSH_DATA_TIME_MS", (int) past);
if (past > 500) {
log.info("Flush data to disk costs {} ms", past);
}
} catch (Throwable e) {
CommitLog.log.warn("{} service has exception. ", this.getServiceName(), e);
this.printFlushProgress();
}
}
// Normal shutdown, to ensure that all the flush before exit
boolean result = false;
for (int i = 0; i < RETRY_TIMES_OVER && !result; i++) {
result = CommitLog.this.mappedFileQueue.flush(0);
CommitLog.log.info("{} service shutdown, retry {} times {}", this.getServiceName(), i + 1, result ? "OK" : "Not OK");
}
...
}
flushCommitLogLeastPages:最少积累页数。如果这次只有一页(4KB)新数据,不值得触发一次 fsync(),等积累到足够的页数再批量刷盘。
ConsumeQueue 异步刷盘
ConsumeQueue 的刷盘比 CommitLog 更宽松——因为它是冗余的索引,丢了可以从 CommitLog 重建。
flushIntervalConsumeQueue = 1000 # ConsumeQueue 刷盘间隔(1s),比 CommitLog 更慢
同步刷盘的可靠性保证:GroupCommit
同步刷盘时,每次消息写入都调用 fsync(),效率很低。RocketMQ 使用 GroupCommit(组提交) 优化:
10 条消息几乎同时到达
│
▼
不分别 fsync(),而是等一批积累后一次 fsync()
│
▼
一次 fsync() 写入 10 条消息 ← 10 倍吞吐量
这引入了微小的延迟(等待一批消息汇聚),但大幅提升了同步刷盘的吞吐量。
🤔 思考穿插:同步刷盘既然每条都要 fsync 落盘,为什么还硬上 GroupCommit 而不是干脆用异步刷盘?—— 因为同步刷盘的诉求是「消息不丢」,GroupCommit 只是把「每条 fsync」改成「一批一次 fsync」,本质没改变「落盘才 ACK」的保证;它用极小的等待延迟换来接近异步的吞吐。所以它是在「可靠性」前提下做性能优化,而不是牺牲可靠性。
// CommitLog.GroupCommitService —— 同步刷盘:把一批请求聚合成一次 fsync
class GroupCommitService extends FlushCommitLogService {
private LinkedList<GroupCommitRequest> requestsWrite = new LinkedList<>();
private LinkedList<GroupCommitRequest> requestsRead = new LinkedList<>();
private final PutMessageSpinLock lock = new PutMessageSpinLock();
public void putRequest(final GroupCommitRequest request) {
lock.lock();
try {
this.requestsWrite.add(request);
} finally {
lock.unlock();
}
this.wakeup();
}
private void swapRequests() {
lock.lock();
try {
LinkedList<GroupCommitRequest> tmp = this.requestsWrite;
this.requestsWrite = this.requestsRead;
this.requestsRead = tmp;
} finally {
lock.unlock();
}
}
private void doCommit() {
if (!this.requestsRead.isEmpty()) {
for (GroupCommitRequest req : this.requestsRead) {
boolean flushOK = CommitLog.this.mappedFileQueue.getFlushedWhere() >= req.getNextOffset();
for (int i = 0; i < 1000 && !flushOK; i++) {
CommitLog.this.mappedFileQueue.flush(0);
flushOK = CommitLog.this.mappedFileQueue.getFlushedWhere() >= req.getNextOffset();
if (flushOK) {
break;
} else {
// When transientStorePoolEnable is true, the messages in writeBuffer may not be committed
// to pageCache very quickly, and flushOk here may almost be false, so we can sleep 1ms to
// wait for the messages to be committed to pageCache.
try {
Thread.sleep(1);
} catch (InterruptedException ignored) {
}
}
}
req.wakeupCustomer(flushOK ? PutMessageStatus.PUT_OK : PutMessageStatus.FLUSH_DISK_TIMEOUT);
}
long storeTimestamp = CommitLog.this.mappedFileQueue.getStoreTimestamp();
if (storeTimestamp > 0) {
CommitLog.this.defaultMessageStore.getStoreCheckpoint().setPhysicMsgTimestamp(storeTimestamp);
}
this.requestsRead = new LinkedList<>();
} else {
// Because of individual messages is set to not sync flush, it
// will come to this process
CommitLog.this.mappedFileQueue.flush(0);
}
}
@Override
public void run() {
CommitLog.log.info("{} service started", this.getServiceName());
while (!this.isStopped()) {
try {
this.waitForRunning(10);
this.doCommit();
} catch (Exception e) {
CommitLog.log.warn("{} service has exception. ", this.getServiceName(), e);
}
}
// Under normal circumstances shutdown, wait for the arrival of the
// request, and then flush
try {
Thread.sleep(10);
} catch (InterruptedException e) {
CommitLog.log.warn("GroupCommitService Exception, ", e);
}
this.swapRequests();
this.doCommit();
CommitLog.log.info("{} service end", this.getServiceName());
}
...
}
分析:生产者线程不亲自 fsync,而是 putRequest 把请求丢进 requestsWrite 后阻塞在 request.future() 上。刷盘线程每 10ms 醒来 doCommit(),用一次 mappedFileQueue.flush(0) 把整批请求的 offset 全部刷到 getFlushedWhere() >= nextOffset,再 wakeupCustomer 统一唤醒——这就是 GroupCommit 组提交,一次 fsync 摊销了一整批消息的落盘开销。
消息确认时机对比
同步刷盘(SYNC_FLUSH):
┌─────────┐ ┌─────────┐ ┌──────────┐ ┌────────┐
│ Producer │────▶│ Broker │────▶│ PageCache │────▶│ Disk │
└─────────┘ └─────────┘ └──────────┘ └────────┘
│ │
│ fsync() │
│◀─────────────│
│
ACK │
◀───┘
ACK 在数据落盘后返回
异步刷盘(ASYNC_FLUSH):
┌─────────┐ ┌─────────┐ ┌──────────┐ ┌────────┐
│ Producer │────▶│ Broker │────▶│ PageCache │─ ─ ▶│ Disk │
└─────────┘ └─────────┘ └──────────┘ └────────┘
│ (延迟写入)
ACK │
◀───┘
ACK 在数据进入 PageCache 后立即返回
生产环境选择建议
| 场景 | 刷盘模式 | 理由 |
|---|---|---|
| 金融支付、交易流水 | SYNC_FLUSH | 消息绝对不能丢失 |
| 订单、通知、日志 | ASYNC_FLUSH | 允许极少情况下的少量丢失 |
| 中间方案 | ASYNC_FLUSH + 主从同步 | 利用主从复制弥补异步刷盘的可靠性 |
| 超高吞吐(>10万 TPS) | ASYNC_FLUSH + SSD | 顺序写 SSD 极限吞吐 |
实践中的折中:大多数公司选择 ASYNC_FLUSH + 同步主从复制(Master-Slave 同步)。即使 Master 宕机导致 Page Cache 数据丢失,Slave 上还有副本。这在不牺牲吞吐量的前提下,提供了接近同步刷盘的可靠性。
总结
RocketMQ 刷盘机制的本质是一个取舍:
同步刷盘 = 可靠性 ↑↑ 吞吐量 ↓↓
异步刷盘 = 可靠性 ↓ 吞吐量 ↑↑
主从复制 = 可靠性 ↑ 吞吐量 ↓(网络开销)
三个关键词:mmap(内存映射文件,绕过用户态缓冲区)、顺序写(Append-only,比随机写快 100 倍)、GroupCommit(批量 fsync,降低同步刷盘的开销)。理解这三者,就理解了 RocketMQ 为什么能在普通磁盘上跑到数十万 TPS。
章末提问
追问 1:同步刷盘和异步刷盘的本质区别是什么?
结论先行:本质区别是 ACK 的返回时机——同步刷盘等 fsync() 真正落盘才返回 ACK,异步刷盘写进 Page Cache 就返回 ACK。
因为:同步刷盘把「落盘」纳入消息确认的前置条件,可靠性最高但吞吐受磁盘顺序写速度限制;异步刷盘把落盘交给后台线程,吞吐高、延迟低,但数据还在 OS 内存缓冲区里,遇到操作系统崩溃会丢。两者的取舍就是「可靠性 vs 吞吐量」。
追问 2:异步刷盘下,Broker 进程崩溃和操作系统崩溃对消息的影响一样吗?
结论先行:不一样——Broker 进程崩溃一般不丢消息,操作系统崩溃会丢 Page Cache 里的数据。
因为:异步刷盘的数据已经进了 Page Cache,这是 OS 管理的共享内存,进程挂了 OS 还在、会继续把脏页刷到磁盘;但如果整个 OS 崩溃(掉电/内核 panic),Page Cache 来不及落盘的部分就永久丢失。所以异步刷盘的「丢」主要针对 OS 级故障,而不是进程级故障。
追问 3:GroupCommit 为什么能提升同步刷盘的吞吐,它牺牲了什么?
结论先行:GroupCommit 把一批并发请求聚合成一次 fsync(),摊薄了单次落盘的系统调用开销,牺牲的是极小的等待延迟。
因为:每条消息单独 fsync() 会让磁盘 IO 次数和系统调用次数成正比,成为吞吐瓶颈;GroupCommit 让刷盘线程每 10ms 醒来一次,用一次 flush 把整批 offset 全部刷到 getFlushedWhere() >= nextOffset,再统一唤醒所有阻塞的请求,一条 fsync 摊销了一整批消息。代价是消息 ACK 要等这批汇聚完成,多出约 10ms 量级的延迟。