You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

大文件量场景下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

额外优化建议

  1. 优雅停止:当前stop()仅设置is_valid_,建议调用CancelIoEx取消未完成的异步请求,确保线程及时退出
  2. 多线程安全:如果后续需要多线程处理事件,需为每个IO请求分配独立的OVERLAPPED和缓冲区对象
  3. 日志细化:增加错误码对应的具体描述,方便排查问题

内容的提问来源于stack exchange,提问作者Kevin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 03:30:55