WatchService 文件监听与 MappedByteBuffer 内存映射源码
概述
WatchService 是 NIO.2 的文件变更监听 API:把目录注册到服务后,目录内文件的创建、修改、删除会以 WatchKey 事件的形式异步推给应用。Windows 上用 IOCP + ReadDirectoryChangesW 实现真正的内核级回调,Linux 上用 inotify,跨平台还有基于轮询的 PollingWatchService 兜底。
MappedByteBuffer 则是内存映射文件:通过 mmap 把文件区域映射到进程地址空间,读写缓冲区直接命中页缓存,零系统调用拷贝。两者都体现了 NIO.2 与操作系统深度协作的设计。本文基于 OpenJDK 21 源码拆解两条链路。
核心源码解析
① FileSystem.newWatchService() 的创建
java
// WindowsFileSystemProvider
@Override
public WatchService newWatchService() throws IOException {
return new WindowsWatchService(this);
}- Windows 上返回
WindowsWatchService(基于 IOCP 实现);Linux 上返回LinuxWatchService(基于 inotify)。 - 构造
WindowsWatchService时会创建一条原生线程(Poller)专门负责从完成端口取事件,JVM 线程不阻塞在系统调用上。 - 目录较少或平台无原生支持时,可用
PollingWatchService:它按固定间隔DirectoryStream扫描目录、对比快照差异来合成事件,代价是轮询开销与秒级延迟。
② WindowsWatchService 的 ReadDirectoryChangesW
c
// Windows native 层(windows_x86_64/libnio)
hIOCP = CreateIoCompletionPort(INVALID_HANDLE_VALUE, NULL, 0, 0); // ① 创建完成端口
pfd->hDir = CreateFile(dir, FILE_LIST_DIRECTORY, // ② 打开目录
FILE_SHARE_READ | FILE_SHARE_WRITE | FILE_SHARE_DELETE,
NULL, OPEN_EXISTING, FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OVERLAPPED, NULL);
CreateIoCompletionPort(pfd->hDir, hIOCP, (ULONG_PTR)completionKey, 0); // ③ 绑定目录到 IOCP
ReadDirectoryChangesW(pfd->hDir, buffer, bufferLen, TRUE, // ④ 异步监视
FILE_NOTIFY_CHANGE_FILE_NAME | FILE_NOTIFY_CHANGE_DIR_NAME
| FILE_NOTIFY_CHANGE_LAST_WRITE,
&bytesReturned, &pfd->overlapped, callback);CreateIoCompletionPort创建/关联完成端口,后续所有异步 IO 完成都会投递到该端口。ReadDirectoryChangesW是异步调用:传入FILE_FLAG_OVERLAPPED,立即返回,目录变化时内核把结果通过 IOCP 完成通知投递回来。Poller线程循环调用GetQueuedCompletionStatus阻塞取完成事件,一个线程管理所有已注册目录——这正是 IOCP 的优势,无需每目录一个线程。
③ Path.register(WatchService watcher, WatchEvent.Kind<?>... events) 的注册
java
// WindowsPath
@Override
public WatchKey register(WatchService watcher, WatchEvent.Kind<?>... events)
throws IOException {
return register(watcher, events, NOFOLLOW_LINKS);
}
@Override
public WatchKey register(WatchService watcher, WatchEvent.Kind<?>[] events,
WatchEvent.Modifier... modifiers) throws IOException {
...
return ((WindowsWatchService) watcher).register(this, events, modifiers);
}- 事件类型
StandardWatchEventKinds只有 4 种:ENTRY_CREATE/ENTRY_MODIFY/ENTRY_DELETE外加OVERFLOW(事件丢失时的兜底信号)。 register内部创建WindowsWatchKey并调用 native 注册:为目录打开句柄、投递一次ReadDirectoryChangesW。一个目录只能注册一次,重复注册复用同一个 key。- 注册成功后立即返回
WatchKey(不阻塞),事件到达后 key 进入 signaling 状态。
④ WatchService.take() 的阻塞等待
java
// AbstractWatchService 的公共骨架
public final WatchKey take() throws InterruptedException {
do {
WatchKey key = queue.take(); // ① 无事件时阻塞
if ((key.signCount() > 0) || key.isValid()) { // ② 校验有效性
return key;
}
} while (true);
}
public final WatchKey poll(long timeout, TimeUnit unit) throws InterruptedException {
WatchKey key = queue.poll(timeout, unit); // ③ 带超时等待
if ((key != null) && ((key.signCount() > 0) || key.isValid())) {
return key;
}
return null;
}AbstractWatchService内部维护LinkedBlockingQueue<WatchKey>:Poller线程收到内核事件后把对应的 key 入队,take()阻塞在queue.take()上。- 出队后还会检查
signCount > 0或 key 有效——若事件已被消费且目录关闭,key 会被丢弃继续循环,避免返回空事件。 poll(timeout)走LinkedBlockingQueue.poll的超时语义;poll()无参版本非阻塞,无事件直接返回null。
⑤ WatchKey.pollEvents() 的事件消费
java
// WindowsWatchKey
public List<WatchEvent<?>> pollEvents() {
synchronized (this) {
List<WatchEvent<?>> result = events; // ① 取出当前事件列表
events = new ArrayList<>(); // ② 清空,下次 poll 重新累积
return result;
}
}
public boolean reset() {
synchronized (this) {
...
if (state == ST_SIGNALLED) {
state = ST_READY; // ③ signaling -> ready
}
return isValid();
}
}- 一个
WatchKey关联一个缓冲区:ReadDirectoryChangesW完成后,Poller解析FILE_NOTIFY_INFORMATION链表,把每个条目转成WindowsWatchEvent(含文件名与事件类型)追加到events。 pollEvents()取走并清空列表;reset()把 key 从 signaling 态恢复为 ready 态并重新投递一次ReadDirectoryChangesW——不 reset 的话该目录不会再产生事件。- 若缓冲区写满来不及消费,native 层返回
OVERFLOW事件,Poller置overflow标志,reset()时会补发一个StandardWatchEventKinds.OVERFLOW。
⑥ MappedByteBuffer.map(FileChannel.MapMode mode, long position, long size)
java
// FileChannelImpl
public MappedByteBuffer map(MapMode mode, long position, long size) throws IOException {
...
long addr = -1;
int mapMode = mode == MapMode.READ_ONLY ? MAP_RO
: mode == MapMode.READ_WRITE ? MAP_RW
: MAP_PV;
try {
addr = map0(mapMode, position, size); // ① native 映射
} catch (OutOfMemoryError x) {
...
}
...
// ② 按地址与大小构造 DirectByteBuffer(UNALIGNED 时可能用分段映射)
int ps = Bits.pageSize();
int sizeInPages = (int)((size + ps - 1) / ps);
int pagePosition = (int)(position - pagePosition(pos));
...
return Util.newMappedByteBufferR(addr + pagePosition, size, this, unmapOnClose);
}map0是 native 方法,返回映射区域的起始虚拟地址;该地址页对齐,因此实际缓冲区从addr + pagePosition开始,大小向上取整到页。- 三种模式映射到不同保护位:
READ_ONLY→PROT_READ、READ_WRITE→PROT_READ|PROT_WRITE、PRIVATE→ 私有写时复制。 - 返回的
DirectByteBuffer视图直接封装裸地址:get/put读写命中映射页,操作系统按需从磁盘加载页缓存。
⑦ map0() 的 mmap 调用
c
// native(mapfile 实现)
JNIEXPORT jlong JNICALL Java_java_nio_MappedByteBuffer_map0(JNIEnv *env, jobject this,
jint prot, jlong off, jlong len) {
void *mapAddress = 0;
jint fd = fdval(env, fdo); // 文件描述符
mapAddress = mmap(mapAddress, // NULL -> 内核选地址
len, // 映射长度
prot, // PROT_READ / PROT_WRITE
MAP_SHARED, // 与其他进程/文件共享页缓存
fd, // 文件句柄
off); // 文件偏移(需页对齐)
if (mapAddress == MAP_FAILED) {
... throw IOException ...
}
return (jlong)cast_address(mapAddress);
}- Linux 调用
mmap(NULL, len, prot, MAP_SHARED, fd, off):off必须页对齐(map0内部已把position规整到页边界)。 MAP_SHARED使映射与文件共享页缓存:缓冲区写入最终会刷回文件,且多进程映射同一文件时可见。- 映射是惰性的:
mmap只建立地址空间映射,不读文件;首次访问页面时触发缺页中断,由内核从磁盘载入。 FileChannelImpl内部会按Bits.maxDirectMemory与maxMappedSize校验映射大小,超限抛异常。
⑧ MappedByteBuffer.force() 的 msync
java
// DirectByteBuffer(MappedByteBuffer 视图)
public final MappedByteBuffer force() {
int ps = Bits.pageSize();
int pagePosition = position - pagePosition(position);
long address = address() + pagePosition;
long a = address - (address % ps);
long b = address + size();
if (b > a + Integer.MAX_VALUE) {
// 大映射分段 flush
}
try {
force0(fd, a, (int)(b - a)); // ① native 落盘
} catch (IOException ioe) {
throw new UncheckedIOException(ioe);
}
return this;
}force0内部调用msync(addr, len, MS_SYNC):把映射页的脏数据同步写回磁盘并等待完成。MS_SYNC是同步模式——方法返回时数据已落盘;MS_ASYNC只调度异步回写立即返回。- 仅
MAP_SHARED映射需要force;PRIVATE映射的写入不会影响文件,force无意义。
⑨ DirectByteBuffer.Deallocator 的 Cleaner 回收
java
// DirectByteBuffer 构造
DirectByteBuffer(long addr, int cap) {
super(-1, 0, cap, cap);
address = addr;
cleaner = Cleaner.create(this, new Deallocator(addr, size)); // ① 注册清理器
}
private static class Deallocator implements Runnable {
public void run() {
if (address == 0) return;
unsafe.freeMemory(address); // ② 释放本机内存
address = 0;
if (unmapper != null) {
unmapper.run(); // ③ 文件映射的 unmap(munmap)
}
}
}Cleaner基于PhantomReference:缓冲区对象不可达时,ReferenceHandler线程执行Deallocator.run()。- 普通直接内存走
unsafe.freeMemory(address)释放;文件映射的缓冲区(MappedByteBuffer)在构造时额外持有unmapper——它内部调用munmap(address, capacity)解除映射。 FileChannelImpl关闭时若映射未回收,会通过unmapOnClose主动触发unmapper.run(),保证映射与底层文件句柄及时释放,避免"删不掉文件"的经典问题。