外观
Stage07|Modern C++ Foundation
Lesson095|综合实践:实现并验证一个小型任务系统
Tags: #C++17 #ThreadPool #CMake #Testing #Stage07Difficulty: ⭐⭐⭐⭐⭐
阶段挑战
本课不是再背一遍 C++ 术语,而是实现一个可停止、可测试的小型任务系统。它必须把以下知识连起来:
text
翻译单元与链接
对象生命周期与 RAII
移动语义与所有权
容器和 lambda
mutex / condition_variable / atomic
CMake 与测试
Sanitizer 与压力验证一、需求和非目标
系统需要:
- 固定数量 worker。
- 接收可调用任务。
- 空闲时休眠,不忙等。
- 多生产者可安全提交。
stop()可重复调用。- 析构前收束线程。
- 任务抛异常不杀死 worker。
- 提供已提交、已完成和失败计数。
本实验不实现优先级、工作窃取、无锁队列和跨进程任务。这些能力会显著扩大正确性边界。
二、项目结构
text
task-system/
├── CMakeLists.txt
├── include/task_system/task_system.hpp
├── src/task_system.cpp
├── app/main.cpp
└── tests/task_system_tests.cpp三、先定义状态机
text
Running
├── submit → 入队并通知
└── stop → Stopping
Stopping
├── worker 继续取已有任务
├── 新 submit 被拒绝
└── 队列为空且任务完成 → Stopped
Stopped
└── stop → 无操作课程选择“排空队列后停止”。状态语义必须先确定,否则锁写得再正确也无法判断提交和关闭的竞态结果。
四、公共接口
cpp
#pragma once
#include <atomic>
#include <condition_variable>
#include <cstddef>
#include <functional>
#include <mutex>
#include <queue>
#include <thread>
#include <vector>
namespace task_system {
class TaskSystem {
public:
explicit TaskSystem(std::size_t workerCount);
~TaskSystem();
TaskSystem(const TaskSystem&) = delete;
TaskSystem& operator=(const TaskSystem&) = delete;
TaskSystem(TaskSystem&&) = delete;
TaskSystem& operator=(TaskSystem&&) = delete;
bool submit(std::function<void()> task);
void stop();
std::size_t submitted() const noexcept;
std::size_t completed() const noexcept;
std::size_t failed() const noexcept;
private:
void workerLoop();
mutable std::mutex mutex_;
std::condition_variable ready_;
std::queue<std::function<void()>> tasks_;
std::vector<std::thread> workers_;
bool stopping_{false};
std::atomic<std::size_t> submitted_{0};
std::atomic<std::size_t> completed_{0};
std::atomic<std::size_t> failed_{0};
};
} // namespace task_system任务系统不能复制或移动,因为 worker 捕获了 this,移动 owner 会让线程保留旧地址。若未来必须可移动,需要额外共享状态对象,而不是简单 = default。
五、核心实现
cpp
#include "task_system/task_system.hpp"
#include <stdexcept>
#include <utility>
namespace task_system {
TaskSystem::TaskSystem(std::size_t workerCount) {
if (workerCount == 0) {
throw std::invalid_argument("workerCount must be greater than zero");
}
workers_.reserve(workerCount);
try {
for (std::size_t i = 0; i < workerCount; ++i) {
workers_.emplace_back([this] { workerLoop(); });
}
} catch (...) {
stop();
throw;
}
}
TaskSystem::~TaskSystem() {
stop();
}
bool TaskSystem::submit(std::function<void()> task) {
if (!task) return false;
{
std::lock_guard<std::mutex> lock(mutex_);
if (stopping_) return false;
tasks_.push(std::move(task));
submitted_.fetch_add(1, std::memory_order_relaxed);
}
ready_.notify_one();
return true;
}
void TaskSystem::workerLoop() {
for (;;) {
std::function<void()> task;
{
std::unique_lock<std::mutex> lock(mutex_);
ready_.wait(lock, [this] {
return stopping_ || !tasks_.empty();
});
if (stopping_ && tasks_.empty()) return;
task = std::move(tasks_.front());
tasks_.pop();
}
try {
task();
completed_.fetch_add(1, std::memory_order_relaxed);
} catch (...) {
failed_.fetch_add(1, std::memory_order_relaxed);
}
}
}
void TaskSystem::stop() {
{
std::lock_guard<std::mutex> lock(mutex_);
stopping_ = true;
}
ready_.notify_all();
for (auto& worker : workers_) {
if (worker.joinable()) worker.join();
}
workers_.clear();
}
std::size_t TaskSystem::submitted() const noexcept {
return submitted_.load(std::memory_order_relaxed);
}
std::size_t TaskSystem::completed() const noexcept {
return completed_.load(std::memory_order_relaxed);
}
std::size_t TaskSystem::failed() const noexcept {
return failed_.load(std::memory_order_relaxed);
}
} // namespace task_system六、逐段解释关键设计
1. 为什么任务在锁外执行
队列锁只保护 stopping_ 和 tasks_。如果在锁内执行任务,一个慢任务会阻止所有 submit 和其他 worker 取任务,还可能因任务重入 submit 而死锁。
2. 为什么条件变量使用谓词
它处理伪唤醒和通知先于等待的情况。worker 只相信受 mutex 保护的“停止或队列非空”状态。
3. 为什么计数器可以 relaxed
三个计数仅用于统计,不用于发布任务结果或同步其他内存。若调用方要等待所有任务完成,不能通过轮询计数替代正确的完成条件变量或 future。
4. 构造失败为什么调用 stop
创建第 N 个线程可能抛出。此前已经启动的线程仍必须被通知和 join,否则异常展开到 std::thread 析构会 terminate。
七、CMake
cmake
cmake_minimum_required(VERSION 3.20)
project(task_system_lab LANGUAGES CXX)
find_package(Threads REQUIRED)
add_library(task_system src/task_system.cpp)
target_include_directories(task_system PUBLIC ${CMAKE_CURRENT_SOURCE_DIR}/include)
target_compile_features(task_system PUBLIC cxx_std_17)
target_link_libraries(task_system PUBLIC Threads::Threads)
add_executable(task_demo app/main.cpp)
target_link_libraries(task_demo PRIVATE task_system)
include(CTest)
if(BUILD_TESTING)
add_executable(task_tests tests/task_system_tests.cpp)
target_link_libraries(task_tests PRIVATE task_system)
add_test(NAME task_tests COMMAND task_tests)
endif()八、最小测试
cpp
#include "task_system/task_system.hpp"
#include <atomic>
#include <cassert>
#include <stdexcept>
int main() {
task_system::TaskSystem system(4);
std::atomic<int> sum{0};
for (int i = 1; i <= 1000; ++i) {
assert(system.submit([&sum, i] { sum.fetch_add(i); }));
}
assert(system.submit([] { throw std::runtime_error("expected"); }));
system.stop();
assert(sum.load() == 1000 * 1001 / 2);
assert(system.submitted() == 1001);
assert(system.completed() == 1000);
assert(system.failed() == 1);
assert(!system.submit([] {}));
system.stop();
}九、必须覆盖的边界测试
text
workerCount=0:明确拒绝
空任务:返回 false
没有任务立即 stop:不挂起
任务执行中 stop:排空后退出
多个 producer 同时 submit:不丢任务
任务内抛异常:worker 继续工作
重复 stop:无崩溃、无重复 join
stop 与 submit 竞争:提交结果符合状态机
任务内调用 stop:当前实现会自 join,必须禁止或重构最后一条非常重要:worker 内调用 stop() 会尝试 join 自己。当前接口应明确要求只能由 owner 线程停止,或者在实现中识别当前线程并设计异步 stop/join 两阶段协议。不要悄悄忽略这个边界。
十、Sanitizer 与压力实验
ASan/UBSan
检查任务捕获和对象析构是否存在越界、悬空或未定义行为。
TSan
8 个 producer 各提交 10000 个任务,同时另一个线程触发 stop,重复数百轮。TSan 应无数据竞态;接受的任务数必须等于 completed + failed。
性能实验
分别测试空任务、100 微秒任务和 10 毫秒任务。记录吞吐、队列峰值和停止延迟。空任务基准主要测调度开销,不代表真实任务收益。
十一、与 Cocos Native 的连接
任务系统只能在后台处理独立 C++ 数据,例如解码、压缩或路径计算。不要在 worker 中直接改 Node、Texture、Renderer 或调用 JavaScript。正确链路是:
text
主线程复制输入并提交
→ worker 计算纯数据结果
→ 投递主线程完成队列
→ 主线程检查 scene/session/owner
→ 更新 Cocos 对象或触发 JS 回调任务系统析构前必须先停止新提交,并确保所有 native 回调不再引用已销毁的绑定对象。
十二、阶段诊断题
1. stop 偶尔永远不返回,应收集什么?
获取所有线程堆栈、队列长度、stopping_ 状态和每个任务开始/结束日志。检查任务是否永久阻塞、是否持锁执行用户代码、是否发生自 join。
2. 计数全部正确,是否证明没有竞态?
不能。某次运行正确不证明 happens-before;仍需审计每个共享状态并使用 TSan 压力测试。
3. 如何支持返回值?
可以用 std::packaged_task 和 std::future 包装任务,但要定义异常传播、取消、future 未消费以及 shutdown 时未执行任务的语义。
十三、Stage07 验收标准
完成者应提交:
- 可从空 build 目录构建的项目。
- 通过的基础和边界测试。
- ASan/UBSan/TSan 的运行记录。
- 一张任务状态机和所有权图。
- 一份压力测试数据,而非“运行没问题”。
- 对 worker 内 stop、自身生命周期和主线程回调边界的说明。
十四、Stage07 总结
Phase 2 的目标不是记住所有 C++ 语法,而是形成可验证的工程链路:
text
用类型表达生命周期和所有权
→ 用容器和算法组织数据
→ 用同步协议保护共享状态
→ 用 CMake 固化构建
→ 用测试、调试器和 Sanitizer 验证完成这一阶段后,才适合继续深入 Cocos 引擎源码、渲染后端、Native 平台能力或更复杂的高并发系统。
