【问题标题】:Strange behaviour of GetQueuedCompletionStatus when used from thread pool worker threads从线程池工作线程中使用 GetQueuedCompletionStatus 时的奇怪行为
【发布时间】:2018-01-15 17:44:03
【问题描述】:

我一直在测试将 IO 完成端口与线程池中的工作线程结合起来,并偶然发现了一种我无法解释的行为。尤其是下面的代码:

  int data;
  for (int i = 0; i < NUM; ++i)
      PostQueuedCompletionStatus(cp, 1, NULL, reinterpret_cast<LPOVERLAPPED>(&data));

  {
      std::thread t([&] ()
      {
            LPOVERLAPPED aux;
            DWORD        cmd;
            ULONG_PTR    key;

            for (int i = 0; i < NUM; ++i)
            {
              if (!GetQueuedCompletionStatus(cp, &cmd, &key, &aux, 0))
                break;
              ++count;
            }
      });

      t.join();
   }

工作得很好,并接收到 NUM 个状态通知(NUM 是大数,100000 或更多),类似的代码使用线程池工作对象,每个工作项读取一个状态通知并在阅读后重新发布工作项,阅读数百个状态通知后失败。具有以下全局变量(请不要介意名称):

HANDLE cport;
PTP_POOL pool;
TP_CALLBACK_ENVIRON env;
PTP_WORK work;
std::size_t num_calls;
std::mutex mutex;
std::condition_variable cv; 
bool job_done;

还有回调函数:

static VOID CALLBACK callback(PTP_CALLBACK_INSTANCE instance_, PVOID pv_, PTP_WORK work_)
{
  LPOVERLAPPED aux;
  DWORD        cmd;
  ULONG_PTR    key;

  if (GetQueuedCompletionStatus(cport, &cmd, &key, &aux, 0))
  {
    ++num_calls;
    SubmitThreadpoolWork(work);
  }
  else
  {
    std::unique_lock<std::mutex> l(mutex);
    std::cout << "No work after " << num_calls << " calls.\n";
    job_done = true;
    cv.notify_one();
  }
}

以下代码:

{
  job_done = false;
  std::unique_lock<std::mutex> l(mutex);

  num_calls = 0;
  cport = CreateIoCompletionPort(INVALID_HANDLE_VALUE, NULL, 0, 1);

  pool = CreateThreadpool(nullptr);
  InitializeThreadpoolEnvironment(&env);
  SetThreadpoolCallbackPool(&env, pool);

  work = CreateThreadpoolWork(callback, nullptr, &env);

  for (int i = 0; i < NUM; ++i)
      PostQueuedCompletionStatus(cport, 1, NULL, reinterpret_cast<LPOVERLAPPED>(&data));

  SubmitThreadpoolWork(work);
  cv.wait_for(l, std::chrono::milliseconds(10000), [] { return job_done; } );
}

尽管 NUM 设置为 1000000,但在调用 GetQueuedCompletionStatus 大约 250 次后会报告“...之后不再工作”。更奇怪的是,将等待时间从 0 设置为 10 毫秒会增加成功呼叫几十万,偶尔会阅读所有 1000000 条通知。我不太明白,因为所有状态通知都是在第一次提交工作对象之前发布的。

将完成端口和线程池结合起来是否真的有问题,或者我的代码有什么问题?请不要谈论我为什么要这样做-我正在调查可能性并偶然发现了这一点。在我看来,它应该可以工作,并且无法弄清楚出了什么问题。谢谢。

【问题讨论】:

  • 您应该检查PostQueuedCompletionStatus(和其他winapi函数)返回的值,如果失败,请检查GetLastError
  • 完整的代码是这样做的,为了简单起见,我删除了检查。未报告任何错误。
  • 你应该把它们加回来。
  • 它们会使示例混乱。此示例报告由 GetQueuedCompletionStatus 返回的不正确(不足)数量的状态通知,无论错误检查如何。特别是,当 GetQueuedCompletionStatus 返回 false 时,它​​会将错误代码设置为 0x102,这表示超时,这反过来又表示没有什么可以返回。没有其他函数报告失败。
  • callback 内部对SubmitThreadpoolWork 的调用似乎是可疑的。这不会导致在尝试修改 num_calls 时在另一个池线程中调用相同的 callback 函数导致竞争条件吗?

标签: c++ windows winapi threadpool


【解决方案1】:

我试过运行这段代码,问题似乎是提供给CreateIoCompletionPortNumberOfConcurrentThreads 参数。传递 1 意味着执行 callback 的第一个池线程与 io 完成端口相关联,但由于线程池可能使用不同的线程执行 callback GetQueuedCompletionStatus 将在这种情况发生时失败。 From documentation:

要仔细考虑的 I/O 完成端口的最重要属性是并发值。完成端口的并发值在通过NumberOfConcurrentThreads 参数使用CreateIoCompletionPort 创建时指定。此值限制与完成端口关联的可运行线程的数量。当与完成端口关联的可运行线程总数达到并发值时,系统会阻止与该完成端口关联的任何后续线程的执行,直到可运行线程数降至并发值以下。

尽管任意数量的线程都可以为指定的 I/O 完成端口调用 GetQueuedCompletionStatus,但当指定的线程第一次调用 GetQueuedCompletionStatus 时,它会与指定的 I/O 完成端口相关联,直到出现以下三种情况之一发生:线程退出,指定不同的I/O完成端口,或者关闭I/O完成端口。换句话说,一个线程最多可以关联一个 I/O 完成端口。

因此,要将 io 完成与线程池一起使用,您需要将并发线程数设置为线程池的大小(您可以使用 SetThreadpoolThreadMaximum 进行设置)。

::DWORD const threads_count{1};

cport = ::CreateIoCompletionPort(INVALID_HANDLE_VALUE, NULL, 0, threads_count);
...
pool = ::CreateThreadpool(nullptr);
::SetThreadpoolThreadMaximum(pool, threads_count);

【讨论】:

    猜你喜欢
    • 2018-09-30
    • 1970-01-01
    • 1970-01-01
    • 2014-03-25
    • 2018-03-06
    • 1970-01-01
    • 2015-01-27
    • 1970-01-01
    相关资源
    最近更新 更多