大文件量场景下ReadDirectoryChangesW+IOCP监控失效问题求助
问题描述
用ReadDirectoryChangesW结合IOCP实现文件夹新增文件监控时,批量复制约800个文件后程序停止工作:
- 调试发现
GetQueuedCompletionStatus返回TRUE,但lpNumberOfBytesTransfered参数为0,后续新增单个文件也无法触发监控 - 触发异常时
OVERLAPPED结构体的internal字段值为268,net helpmsg查不到对应错误信息
代码实现
Windows API封装类(FileIOAdapterImpl)
//FileIOAdapterImpl HandlePtr FileIoAdapterImpl::CreateIoCompletionPortWrapper() { HandlePtr io_handle{ CreateIoCompletionPort(INVALID_HANDLE_VALUE, nullptr, 0, 0) }; return io_handle; } bool FileIoAdapterImpl::AssociateWithFileHandle(HANDLE directory, HANDLE existing_completion_port) { return CreateIoCompletionPort(directory, existing_completion_port, (ULONG_PTR)directory, 1); } bool FileIoAdapterImpl::GetQueuedCompletionStatusWrapper(HANDLE completion_token, LPDWORD bytes_transferred, PULONG_PTR completion_key, LPOVERLAPPED* overlap) { return GetQueuedCompletionStatus(completion_token, bytes_transferred, completion_key, overlap, 16); } HandlePtr FileIoAdapterImpl::CreateFileWrapper(const std::string& path) { HandlePtr dir_handle{ CreateFileA(path.c_str(), FILE_LIST_DIRECTORY, FILE_SHARE_READ | FILE_SHARE_WRITE | FILE_SHARE_DELETE, NULL, OPEN_EXISTING, FILE_FLAG_BACKUP_SEMANTICS | FILE_FLAG_OVERLAPPED, NULL) }; return dir_handle; } bool FileIoAdapterImpl::ReadDirectoryChangesWrapper(HANDLE directory, std::vector<std::byte>& buffer, LPDWORD bytes_returned, OVERLAPPED* overlapped) { return ReadDirectoryChangesW(directory, buffer.data(), buffer.size(), FALSE, FILE_NOTIFY_CHANGE_FILE_NAME, bytes_returned, overlapped, nullptr); }
注:HandlePtr是带自定义删除器的unique_ptr,自动调用CloseHandle
ModernDirectoryWatcher类定义(头文件)
//ModernDirectoryWatcher.h class ModernDirectoryWatcher { public: using Callback = std::function<void(const std::string& filename)>; ModernDirectoryWatcher(const std::filesystem::path& input_directory, std::vector<std::string> file_types, Callback callback, FileIOAdapterPtr io_wrapper = std::make_shared<FileIoAdapterImpl>()); explicit operator bool() const { return is_valid_; } bool watch(); void stop(); private: bool event_recv(); bool event_send(); void handle_events(); bool has_event() const { return event_buf_len_ready_ != 0; } bool is_processable_file(const std::filesystem::path& path) const; Poco::Logger& logger_{ Poco::Logger::get("ModernDirectoryWatcher") }; std::filesystem::path input_directory_; std::vector<std::string> file_types_{}; Callback callback_; FileIOAdapterPtr win_io_api_wrapper_; HandlePtr path_handle_; HandlePtr event_completion_token_; unsigned long event_buf_len_ready_{ 0 }; bool is_valid_{ false }; OVERLAPPED event_overlap_{}; std::vector<std::byte> event_buf_{ 64 * 1024 }; };
ModernDirectoryWatcher实现(cpp文件)
//ModernDirectoryWatcher.cpp ModernDirectoryWatcher::ModernDirectoryWatcher(const std::filesystem::path& input_directory, std::vector<std::string> file_types, Callback callback, FileIOAdapterPtr io_wrapper) : input_directory_(input_directory), file_types_(std::move(file_types)), callback_(std::move(callback)), win_io_api_wrapper_(std::move(io_wrapper)) { path_handle_ = win_io_api_wrapper_->CreateFileWrapper(input_directory_.string()); if (path_handle_.get() != INVALID_HANDLE_VALUE) { poco_information(logger_, "Create Completition Token"); event_completion_token_ = win_io_api_wrapper_->CreateIoCompletionPortWrapper(); } if (event_completion_token_.get() != nullptr) { poco_information(logger_, "Associate with the FileHandle"); is_valid_ = win_io_api_wrapper_->AssociateWithFileHandle(path_handle_.get(), event_completion_token_.get()); } } bool ModernDirectoryWatcher::watch() { poco_information(logger_, "Start Watching..."); if (is_valid_) { poco_information(logger_, "Receive Events"); event_recv(); while (is_valid_ && has_event()) { poco_information(logger_, "There are still events to process..."); event_send(); } while (is_valid_) { ULONG_PTR completion_key{ 0 }; LPOVERLAPPED overlap{ 0 }; bool complete = win_io_api_wrapper_->GetQueuedCompletionStatusWrapper(event_completion_token_.get(), &event_buf_len_ready_, &completion_key, &overlap); if (complete && event_buf_len_ready_ == 0) { poco_error(logger_, "Error"); // HERE THE ERROR IS PRINTED } if (complete && overlap) { poco_information(logger_, "Handle the events"); handle_events(); } else if (int err_code = GetLastError() != 258 && !complete) { poco_error(logger_, "Error"); } } return true; } else { return false; } } void ModernDirectoryWatcher::handle_events() { while (is_valid_ && has_event()) { poco_information(logger_, "Send Event"); event_send(); poco_information(logger_, "Receive Events"); event_recv(); } } void ModernDirectoryWatcher::stop() { poco_notice(logger_, "Stop the Watcher"); is_valid_ = false; } bool ModernDirectoryWatcher::event_recv() { event_buf_len_ready_ = 0; DWORD bytes_returned = 0; memset(&event_overlap_, 0, sizeof(OVERLAPPED)); poco_information(logger_, "Call the ReadDirectoryChanges"); auto read_ok = win_io_api_wrapper_->ReadDirectoryChangesWrapper(path_handle_.get(), event_buf_, &bytes_returned, &event_overlap_); if (!event_buf_.empty() && read_ok) { event_buf_len_ready_ = bytes_returned > 0 ? bytes_returned : 0; poco_information_f1(logger_, "Event Buffer Len: %?d", event_buf_len_ready_); return true; } if (GetLastError() == ERROR_IO_PENDING) { poco_error(logger_, "Error Pending IO Received stopping..."); event_buf_len_ready_ = 0; is_valid_ = false; } else { poco_error_f1(logger_, "Error Code: %?d", GetLastError()); } return false; } bool ModernDirectoryWatcher::event_send() { auto buf = reinterpret_cast<FILE_NOTIFY_INFORMATION*>(event_buf_.data()); if (is_valid_) { while (buf + sizeof(FILE_NOTIFY_INFORMATION) <= buf + event_buf_len_ready_) { poco_information(logger_, "Get valid Buffer..."); auto filename = input_directory_ / std::wstring{ buf->FileName, buf->FileNameLength / 2 }; if ((buf->Action == FILE_ACTION_ADDED || buf->Action == FILE_ACTION_RENAMED_NEW_NAME) && is_processable_file(filename)) { callback_(filename.string()); } if (buf->NextEntryOffset == 0) { break; } buf = reinterpret_cast<FILE_NOTIFY_INFORMATION*>(reinterpret_cast<std::byte*>(buf) + buf->NextEntryOffset); } return true; } else { return false; } } bool ModernDirectoryWatcher::is_processable_file(const std::filesystem::path& path) const { std::string extension = path.extension().string().erase(0, 1); extension = Poco::toLower(extension); return std::find(file_types_.begin(), file_types_.end(), extension) != file_types_.end(); }
调用方式
ModernDirectoryWatcher directory_watcher{dir, ftypes, [this](const std::string& filename) {...}}; std::thread watcher_thread = std::thread{&ModernDirectoryWatcher::watch, &directory_watcher};
问题分析与修复方案
核心问题1:错误处理ERROR_IO_PENDING
ERROR_IO_PENDING是异步IO的正常状态,表示请求已提交等待完成,但当前代码直接终止监控,这是批量复制时程序崩溃的主要原因。
修复代码(event_recv函数):
bool ModernDirectoryWatcher::event_recv() { event_buf_len_ready_ = 0; DWORD bytes_returned = 0; memset(&event_overlap_, 0, sizeof(OVERLAPPED)); poco_information(logger_, "Call the ReadDirectoryChanges"); auto read_ok = win_io_api_wrapper_->ReadDirectoryChangesWrapper(path_handle_.get(), event_buf_, &bytes_returned, &event_overlap_); if (read_ok) { event_buf_len_ready_ = bytes_returned; poco_information_f1(logger_, "Event Buffer Len: %?d", event_buf_len_ready_); return true; } DWORD err_code = GetLastError(); if (err_code == ERROR_IO_PENDING) { // 异步请求已提交,等待IOCP通知,属于正常流程 poco_information(logger_, "IO request pending, waiting for completion"); event_buf_len_ready_ = 0; return true; } else { poco_error_f1(logger_, "ReadDirectoryChanges failed with error: %?d", err_code); is_valid_ = false; return false; } }
核心问题2:零字节完成未重新投递请求
GetQueuedCompletionStatus返回TRUE但字节数为0时,当前代码仅打印错误,未重新提交ReadDirectoryChangesW请求,导致后续无监控事件。
修复代码(watch函数):
bool ModernDirectoryWatcher::watch() { poco_information(logger_, "Start Watching..."); if (is_valid_) { poco_information(logger_, "Initialize first IO request"); if (!event_recv()) { poco_error(logger_, "Failed to submit initial ReadDirectoryChanges request"); is_valid_ = false; return false; } while (is_valid_) { ULONG_PTR completion_key{ 0 }; LPOVERLAPPED overlap{ 0 }; DWORD bytes_transferred = 0; // 改用INFINITE超时,减少不必要的线程唤醒 bool complete = win_io_api_wrapper_->GetQueuedCompletionStatusWrapper(event_completion_token_.get(), &bytes_transferred, &completion_key, &overlap, INFINITE); if (!complete) { DWORD err_code = GetLastError(); if (err_code != WAIT_TIMEOUT) { poco_error_f1(logger_, "GetQueuedCompletionStatus failed with error: %?d", err_code); is_valid_ = false; } continue; } event_buf_len_ready_ = bytes_transferred; if (bytes_transferred == 0) { poco_warning(logger_, "Zero bytes transferred, re-submitting ReadDirectoryChanges request"); // 重新投递请求,恢复监控 if (!event_recv()) { is_valid_ = false; } continue; } if (overlap) { poco_information(logger_, "Handle the events"); handle_events(); } } return true; } else { return false; } }
核心问题3:缓冲区遍历边界错误
原代码的缓冲区越界判断逻辑错误,可能导致读取垃圾数据。
修复代码(event_send函数):
bool ModernDirectoryWatcher::event_send() { std::byte* buf_end = event_buf_.data() + event_buf_len_ready_; auto buf = reinterpret_cast<FILE_NOTIFY_INFORMATION*>(event_buf_.data()); if (is_valid_ && event_buf_len_ready_ > 0) { while (reinterpret_cast<std::byte*>(buf) + sizeof(FILE_NOTIFY_INFORMATION) <= buf_end) { poco_information(logger_, "Get valid Buffer..."); // 用sizeof(wchar_t)计算文件名长度,避免编码错误 auto filename = input_directory_ / std::wstring{ buf->FileName, buf->FileNameLength / sizeof(wchar_t) }; if ((buf->Action == FILE_ACTION_ADDED || buf->Action == FILE_ACTION_RENAMED_NEW_NAME) && is_processable_file(filename)) { callback_(filename.string()); } if (buf->NextEntryOffset == 0) { break; } std::byte* next_buf = reinterpret_cast<std::byte*>(buf) + buf->NextEntryOffset; // 检查下一个条目是否在缓冲区范围内 if (next_buf >= buf_end) { poco_warning(logger_, "Truncated file notify information, partial data discarded"); break; } buf = reinterpret_cast<FILE_NOTIFY_INFORMATION*>(next_buf); } return true; } else { return false; } }
关于OVERLAPPED.internal=268的说明
错误码268对应ERROR_NO_MORE_FILES,通常是批量事件触发导致缓冲区溢出、事件丢失引起的。可以增大监控缓冲区降低概率:
std::vector<std::byte> event_buf_{ 256 * 1024 }; // 从64k改为256k
额外优化建议
- 优雅停止:当前
stop()仅设置is_valid_,建议调用CancelIoEx取消未完成的异步请求,确保线程及时退出 - 多线程安全:如果后续需要多线程处理事件,需为每个IO请求分配独立的OVERLAPPED和缓冲区对象
- 日志细化:增加错误码对应的具体描述,方便排查问题
内容的提问来源于stack exchange,提问作者Kevin
相关产品推荐
相关产品推荐

