一.最简单方案为什么不够一个最简单方案submit(task)↓ 直接run_process(task)如果这个任务要执行十秒那么这十秒就无法做其他事所以一种自然想法是每来一个任务就新建一个线程去执行。例如任务1→ 线程1任务2→ 线程2任务3→ 线程3...但是当有10000个任务时线程本身也需要资源而且过多线程还会带来大量调度开销。所以采用fixed worker pool 固定工作线程池 bounded FIFO queue 有界先进先出队列worker pool固定4 个 worker可以先理解成worker1worker2worker3worker4这些 worker 会不断从任务队列里取任务 ↓ 执行 ↓ 执行完再回来取下一个所以即使来了 100 个任务100个任务 ↓ 先排队 ↓ 最多4个 worker 同时取任务执行不会创建 100 个线程。TaskForge 主进程 │ ├── worker 线程 ← 线程进程内部 │ │ │ └──fork()│ │ │ └── child 进程 ← 进程worker 造出来的新进程 │ └── 其他 worker 线程bounded FIFO queueFIFO First In, First Out先进先出而 bounded 的意思是队列长度不是无限的有最大容量。比如queue_capacity10最多允许 10 个“还在等待、还没被 worker 取走”的任务留在队列里。queue_capacity 只限制等待队列里的任务数量不把已经被 worker 取走、正在执行的任务算进去。所以4个任务正在执行10个任务在队列里等是被允许的~~ ~~~~ ~~~~ ~~~~ ~~~~ ~~~~ ~~~~ ~~~~ ~~~~ ~~~~ ~~~~ ~~~~Q如果有 4 个 workerqueue_capacity 10当前 4 个任务正在执行、10 个任务正在队列里等待这时候再来第 15 个任务它还能正常进入等待队列吗A不会等待队列腾位置而是直接返回 queue_full。不是继续等待因为submit是nonblocking submit非阻塞提交这就是 backpressure背压系统接不下了就明确告诉上游“我满了”而不是自己无限堆积任务。调用方可以自己决定稍后要不要重试。mutex假设现在有一个共享任务队列queue:[A][B][C]可能同时有线程1submit新任务要 push 进去 线程2submit另一个任务也要 push worker要从队头 pop 一个任务如果三个线程同时直接改这个 queue就可能把内部状态搞乱。所以要用mutex互斥锁保证同一时刻只有一个线程进入操作这份共享数据的关键区域。线程 A ↓ 拿到 mutex ↓ 操作 queue ↓ 释放 mutex 线程 B ↓ 之前必须等 ↓ A 释放后才能操作 queue而critical section临界区就是“拿着这把锁、正在访问共享 queue 的那段代码。”condition_variable假设一个 worker 现在没任务可做。最笨的写法可能是while(queue.empty()){// 一直检查}这叫不断轮询。但是queue 还是空 → worker 检查一次 还是空 → 又检查一次 还是空 → 再检查一次会白白浪费cpu所以更好的方式是没任务时让 worker 睡眠有新任务时再把它叫醒。但是worker 不是“被叫醒了就直接拿任务”而是被唤醒 ↓ 重新拿到 mutex ↓ 重新检查条件//也就是wait(lock, predicate)其中predicate条件谓词 “现在到底能不能继续”的判断条件↓ 条件真的满足才继续为什么醒了还要重新检查因为存在spurious wakeup虚假唤醒也就是worker 有可能醒了但实际上并没有新任务。所以不能写成 收到通知 → 肯定有任务 而应该理解成 收到通知 → “你可以再检查一下了”通知只是提醒重新看看不是发一张“保证有任务”的票但是不能worker睡觉时仍拿着锁因为submit 想拿 queue_mutex ↓ 拿不到 ↓ 没法把新任务 push 进 queue正确路径worker 拿锁 ↓ 检查 queue ↓ 发现没有任务 ↓ condition_variable::wait ↓ 自动释放 mutex ↓ worker 睡眠总流程没活 → wait → 释放锁 → 睡眠 新任务到来 → submit 拿锁 → push → notify worker 醒来 → 重新拿锁 → 重新检查 predicatesubmit队列满了 → 不在这里等 → 直接 queue_full 队列没满 → 接收任务 → 很快返回 accepted 但任务真正完成 → 以后再拿结果accepted只表示执行器已经接管了这个任务。它不代表任务已经开始执行更不代表exec 已成功Q为什么 TaskForge 不能让 submit() 直接返回最终 ProcessResult而要搞一个“以后再取结果”的 future/handle 机制A因为 submit() 的目标就是把任务交给执行器然后尽快返回而不是等任务真正跑完。如果 submit() 要直接返回最终的 ProcessResult那它就只能这样submit(task)↓ 等 worker 拿到任务 ↓ 等run_process()↓ 等 child 执行结束 ↓ 拿到 ProcessResult ↓ submit 才返回会导致调用 submit() 的那个线程会被任务执行时间卡住任务执行完后调用方去哪里拿结果promise futurepromise生产结果的一端 future等待/获取结果的一端promise 是 worker 以后写结果的位置future 是调用方等待结果的通道两者连接到同一个 shared state共享状态worker 调用方 │ │ │ 把结果放进去 │ 等结果 ▼ ▼ promise ─── shared state ─── future接submit的question第一步submit(task)→ “任务我接收了” → 很快返回 第二步 worker 以后执行任务 → 得到 ProcessResult 第三步 调用方以后通过 future/handle → 获取最终结果future 和 shared_future的区别普通std::future通常只适合一次 get()而 std::shared_future 可以让多个持有者等待/获取同一个结果所以更适合可复制的 TaskHandle排队取消与 shutdown 的边界排队取消假设 C 还在 queue 里这时有人请求cancel CTaskForge 不会立刻从队列中间把 C 删除掉。它仍然可能暂时占着那个 queue slot直到 worker 后面把 C 取出来发现它已经被取消然后直接完成 cancellation不会 fork workloadshutdown他不是shutdown()→ 立刻把所有任务都取消而是shutdown()↓ 先停止接收新任务 ↓ 已经 accepted 的任务继续跑完 ↓ 最后等待 worker 全部退出为什么调用 join() 时不能一直拿着 queue_mutexjoin() 是线程里的一个基本操作可以先把它理解成等待某个线程执行结束。std::threadworker(...);worker.join();这里 join() 的意思不是“把线程关掉”而是当前线程停在这里等 worker 自己执行完退出。假设 shutdown 线程这样做shutdown 线程 ↓ 拿到 queue_mutex ↓ 开始join(worker)↓ 等待 worker 结束但 worker 要结束前可能还需要worker ↓ 拿 queue_mutex ↓ 检查 queue/accepting ↓ 确认没任务了 ↓ 退出于是shutdown 拿着 queue_mutex → 等 worker 退出 worker 想拿 queue_mutex → 拿不到 → 无法走到退出就可能形成死锁
阅读完成 · 觉得有帮助?