MiddleКодЧастоЕщё не отвечали
Напишите потокобезопасный пул потоков.
Реализуйте пул потоков, который заранее создаёт N рабочих потоков и выполняет задачи, поданные в общую очередь.
Требования:
- очередь задач — общее состояние; защитите её
std::mutexиstd::condition_variable - рабочие должны ждать с предикатом (
wait(lock, pred)), чтобы ложные пробуждения не извлекали из пустой очереди submit()кладётstd::function<void()>и будит ровно один поток (notify_one)- деструктор ставит флаг остановки, будит всех рабочих, даёт им опустошить очередь и делает
joinкаждого потока (неdetach)
class ThreadPool {
public:
explicit ThreadPool(std::size_t threads);
~ThreadPool();
ThreadPool(const ThreadPool&) = delete;
ThreadPool& operator=(const ThreadPool&) = delete;
void submit(std::function<void()> task);
// ваш код здесь
};
Допишите реализацию.
Заранее создайте N рабочих, ожидающих общей очереди задач, защищённой мьютексом и condition variable. submit() кладёт std::function<void()> и будит один поток; деструктор ставит флаг остановки и делает join.
- ✗Не перепроверять предикат после пробуждения
wait— ложные пробуждения доставляют поток с пустой очередью; всегда используйтеwait(lock, pred)или цикл - ✗Уничтожать пул при наличии задач в очереди — рабочие должны опустошить очередь перед выходом, или деструктор должен решить, отменять или завершить незавершённую работу
- ✗Блокировать
submit()при полной очереди без ограничения размера — неограниченные очереди могут исчерпать память под нагрузкой
- →Как добавить поддержку приоритетов в очередь задач?
- →Как вернуть
std::futureизsubmit(), чтобы вызывающие могли ждать результатов?
Оглавление
Потокобезопасный пул потоков
Реализация
#include <condition_variable>
#include <functional>
#include <mutex>
#include <queue>
#include <thread>
#include <vector>
class ThreadPool {
public:
explicit ThreadPool(std::size_t threads) {
for (std::size_t i = 0; i < threads; ++i)
workers_.emplace_back([this] { workerLoop(); });
}
~ThreadPool() {
{
std::lock_guard lock(mu_);
stop_ = true;
}
cv_.notify_all();
for (auto& t : workers_)
t.join();
}
// Запрет копирования/перемещения
ThreadPool(const ThreadPool&) = delete;
ThreadPool& operator=(const ThreadPool&) = delete;
void submit(std::function<void()> task) {
{
std::lock_guard lock(mu_);
if (stop_)
throw std::runtime_error("submit on stopped ThreadPool");
tasks_.push(std::move(task));
}
cv_.notify_one();
}
private:
void workerLoop() {
while (true) {
std::function<void()> task;
{
std::unique_lock lock(mu_);
cv_.wait(lock, [this] { return stop_ || !tasks_.empty(); });
if (stop_ && tasks_.empty())
return;
task = std::move(tasks_.front());
tasks_.pop();
}
task();
}
}
std::vector<std::thread> workers_;
std::queue<std::function<void()>> tasks_;
std::mutex mu_;
std::condition_variable cv_;
bool stop_ = false;
};
Расширение: submit с std::future
#include <future>
#include <type_traits>
template<typename F>
auto ThreadPool::submitFuture(F&& f) -> std::future<std::invoke_result_t<F>> {
using R = std::invoke_result_t<F>;
auto task = std::make_shared<std::packaged_task<R()>>(std::forward<F>(f));
std::future<R> fut = task->get_future();
submit([task] { (*task)(); });
return fut;
}
Тест
#include <atomic>
#include <cassert>
#include <chrono>
int main() {
ThreadPool pool(4);
std::atomic<int> counter{0};
const int N = 100;
for (int i = 0; i < N; ++i)
pool.submit([&counter] { ++counter; });
// Деструктор пула дожидается завершения всех задач
// (пул дренирует очередь перед join)
}
// После деструктора: counter == N (100)
Ключевые точки
| Аспект | Решение |
|---|---|
| Синхронизация очереди | std::mutex + std::condition_variable |
| Безопасная остановка | флаг stop_ + notify_all() в деструкторе |
| Дренаж очереди | рабочий выходит только когда stop_ && tasks_.empty() |
| Ложные пробуждения | cv_.wait(lock, predicate) |
| Исключения в задачах | передаются через std::packaged_task при использовании submitFuture |
Оглавление