|
|
f1 [! {, D; G5 B0 D7 J
<h1 id="如何实现一个线程池">如何实现一个线程池</h1>6 _, v1 C- G/ K$ E4 `. D
<p>线程池:一种线程使用模式。线程过多会带来调度开销,进而影响缓存局部性和整体性能。而线程池维护着多个线程,等待着监督管理者分配可并发执行的任务。这避免了在处理短时间任务时创建与销毁线程的代价。线程池不仅能够保证内核的充分利用,还能防止过分调度。可用线程数量应该取决于可用的并发处理器、处理器内核、内存、网络sockets等的数量。 例如,对于计算密集型任务,线程数一般取cpu数量+2比较合适,线程数过多会导致额外的线程切换开销。</p>, K' x3 O! S* L) [
<p>如何定义线程池Pool呢,首先最大线程数量肯定要作为线程池的一个属性,并且在new Pool时创建指定的线程。</p>
& w; n* P9 ~6 _: s+ M f) D8 E* X<p>线程池Pool</p>
" E t: t* P1 L4 I<pre><code>pub struct Pool {* U: b$ z r% X8 ?3 x+ G0 X4 t; Z6 U
max_workers: usize, // 定义最大线程数
2 q6 G/ k6 B* a, \}, _7 D" R3 L8 J5 U9 X5 `0 P' h6 ~
; `0 g. _$ Z7 C/ ~impl Pool {
1 X' j, ^; Q6 ]$ y6 f; }" i k fn new(max_workers: usize) -> Pool {}# n# r4 ]" x! A3 H
fn execute<F>(&self, f:F) where F: FnOnce() + 'static + Send {}
) a2 Z" f% R8 I) o6 _$ v}
9 ]6 `' e: u1 U
; |6 E3 b3 P# N: \0 s$ {" Y2 S</code></pre>
% e+ y7 I z i6 x+ a2 y. z<p>用<code>execute</code>来执行任务,<code>F: FnOnce() + 'static + Send</code> 是使用thread::spawn线程执行需要满足的trait, 代表F是一个能在线程里执行的闭包函数。</p>
5 W \4 ~0 a- b T4 H4 L, p- d<p>另一点自然而然会想到在Pool添加一个线程数组, 这个线程数组就是用来执行任务的。比如<code>Vec<Thread></code> balabala。这里的线程是活的,是一个个不断接受任务然后执行的实体。<br>+ }- A- g7 x& o$ L( c) P3 u+ O
可以看作在一个线程里不断执行获取任务并执行的Worker。</p>/ H4 |3 j' \( D" i x
<pre><code>struct Worker where
. g! W3 u; b) O6 B{
$ Z. m' [& {* M! y6 z$ I _id: usize, // worker 编号
2 b5 a2 q& N" |- q& p+ E+ Q}7 z! `% p9 g& A' H( \9 |* p
</code></pre>+ e7 E& v+ o) \8 u+ M; R
<p>要怎么把任务发送给Worker执行呢?mpsc(multi producer single consumer) 多生产者单消费者可以满足我们的需求,<code>let (tx, rx) = mpsc::channel()</code> 可以获取到一对发送端和接收端。<br>3 O$ w/ I6 \) M
把发送端添加到Pool里面,把接收端添加到Worker里面。Pool通过channel将任务发送给多个worker消费执行。</p>
' w; s' u' B2 I# S$ Y/ Z<p><strong>这里有一点需要特别注意,channel的接收端receiver需要安全的在多个线程间共享</strong>,因此需要用<code>Arc<Mutex::<T>></code>来包裹起来,也就是用锁来解决并发冲突。</p>
^, e, k2 F' Y `3 Q<p>Pool的完整定义</p># s4 t/ B7 ~, P1 `) @8 S: {
<pre><code>pub struct Pool {+ Q) y( B9 g3 R
workers: Vec<Worker>,
. [; @. x5 q6 N6 M7 e3 E% ~7 W! [ max_workers: usize,+ U$ [' a& n* {" h
sender: mpsc::Sender<Message>
$ g c5 e; O5 d* w |' T+ t( T}
2 Z' @7 r. P' j. z</code></pre>' g1 F9 g! Q w) u
<p>该是时候定义我们要发给Worker的消息Message了<br>" k# A* e$ @# M% \
定义如下的枚举值</p>6 ?5 R( h! k& a E4 S
<pre><code>type Job = Box<dyn FnOnce() + 'static + Send>;) N- q7 D+ D' w8 ~
enum Message {$ m7 p6 ]' {( S8 L
ByeBye,
5 ?5 g3 m5 O7 M: R7 [; {- \/ m NewJob(Job),- Z/ ~- M/ D: Z
}
' s# { l2 _$ z1 I; x+ ^</code></pre>- ^% \7 q G* q L' N1 q9 [5 s( ~) n
<p>Job是一个要发送给Worker执行的闭包函数,这里ByeBye用来通知Worker可以终止当前的执行,退出线程。</p>
* k- s# L) F% Z. T( H<p>只剩下实现Worker和Pool的具体逻辑了。</p>4 D1 }. K4 b/ f( |. ]3 w, f
<p>Worker的实现</p>
5 k! I0 v" t0 {7 q: F) R( |! p<pre><code>impl Worker
8 P) G: o6 \7 P( Z" X9 v{+ g8 B7 R+ l' Y" n
fn new(id: usize, receiver: Arc::<Mutex<mpsc::Receiver<Message>>>) -> Worker {$ L3 t1 n2 h, ~
let t = thread::spawn( move || {% [' ?( p$ U9 o. Z; Z
loop {
8 u. K/ `6 }% {6 X E* k let receiver = receiver.lock().unwrap();
2 S4 W* \& _* z: }( Z: | let message= receiver.recv().unwrap();
! D# r% X6 \8 l- t match message {
* Y/ d7 F- g- j) K8 O Message::NewJob(job) => {
6 ~4 ? s1 K/ m1 D! O println!("do job from worker[{}]", id);$ t" D9 ~7 B* X) [
job();
# F# N- V/ w* N# i! t `. w3 a },5 S4 @& ]8 s7 y4 P5 x+ Z9 U- S/ f
Message::ByeBye => {
% ]- E5 L9 m/ I! g% e' k4 Z% Q2 j println!("ByeBye from worker[{}]", id);
. n& u- n& {' E break
+ r' Z' X! F4 g7 P1 o3 r },
3 S0 \" x; H7 B5 m" Y8 C1 @; b } / @) a5 s' R2 d0 W
}9 z' K) I: `9 m8 \- o4 F
});, w5 U! \6 E, u( [9 s7 Z
, x8 T# Q& \( N7 u Worker {9 D+ v2 m6 ]) n1 A; T1 [* i
_id: id,
! l! |# h3 f x2 M) j; T! D- O t: Some(t),8 {) n) V) F/ e- {% U1 e7 Z
}# c' ?' b8 n( H+ P4 g5 O
}
+ i$ {# P d5 |9 {}, O; I m7 d" \: p4 f$ F# C$ p2 v) W
</code></pre> V+ n) X& T) b0 D# D' k4 x
<p><strong>let message = receiver.lock().unwrap().recv().unwrap();</strong> 这里获取锁后从receiver获取到消息体,然后let message结束后rust的生命周期会自动释放掉锁。<br>
9 ?: k8 m) `8 I* V# @" H但如果写成</p>
7 u6 k8 [. z- f) U<pre><code>while let message = receiver.lock().unwrap().recv().unwrap() {' @. C! z/ W# |0 Q
};
% }3 P' N- o( {# \</code></pre>$ n0 m/ u7 S. V
<p>while let 后面整个括号都是一个作用域,要在这个作用域结束后,锁才会释放,比上面let message要锁定久时间。<br># N" ~) ~3 u! p) {' B; W" ~
rust的mutex锁没有对应的unlock方法,由mutex的生命周期管理。</p>$ n* c0 k3 c/ D5 H! V0 k. m! R
<p>我们给Pool实现<code>Drop</code> trait, 让Pool被销毁时,自动暂停掉worker线程的执行。</p>
- _( y% v# ~ g; \<pre><code>impl Drop for Pool {
7 ~; o' ^, w7 m3 f fn drop(&mut self) {6 e* E9 c4 z: K6 m$ i' e' z
for _ in 0..self.max_workers {
t( ^3 r1 y+ |2 [1 M0 ?5 v5 e self.sender.send(Message::ByeBye).unwrap();' R6 \" y4 U1 T
}) p# X9 w0 p! m
for w in self.workers.iter_mut() {, Z. z2 Y9 e% T/ t+ y8 l- b
if let Some(t) = w.t.take() {
: T# X' f% f6 k+ i* c1 J t.join().unwrap();
9 T! r9 b# V- y% b* J }; O0 }6 y0 P7 z8 e- {1 x3 I
}. g( [8 P! a' e& J; p6 S4 W
}
/ y; q4 T- |. Z: Y) U9 f- W} g5 F/ n2 O, l& l3 h
v/ O; o H$ I9 s</code></pre>& n$ o, ?6 N5 ?
<p><strong>drop方法里面用了两个循环</strong>,而不是在一个循环里做完两件事?</p>$ A0 I: a( j! g6 H# v# A- C
<pre><code>for w in self.workers.iter_mut() {& a2 i9 W# G7 F! \, W
if let Some(t) = w.t.take() {
. j8 `3 O( R5 ?* U self.sender.send(Message::ByeBye).unwrap();
* t7 |7 \# D% e. I8 m t.join().unwrap();
$ M' k/ k* K6 J$ ~ }
/ H# z( u& F! t Z2 j}
( \# _; X/ _& G7 A" r& p' x+ \ i! h3 ^7 f8 C5 R+ y9 J3 [, ]% `
</code></pre>
8 ?- G" Q+ f5 G7 B* t+ c3 ]# K7 E<p>这里面隐藏了一个会造成死锁的陷阱,比如两个Worker, 在单个循环里面迭代所有Worker,再将终止信息发送给通道后,直接调用join,<br>( A) j) z O: ~% p- `9 S5 C# e
我们预期是第一个worker要收到消息,并且等他执行完。当情况可能是第二个worker获取到了消息,第一个worker没有获取到,那接下来的join就会阻塞造成死锁。</p>. H% x/ _9 P6 h: a" S9 B" a
<p><strong>注意到没有,Worker是被包装在Option内的</strong>,这里有两个点需要注意</p> [5 k! g/ u+ |
<ol>
! q% ]4 a8 U) l/ ~4 Y, I<li>t.join 需要持有t的所有权</li>
7 H9 B: l. U+ W7 b% ^2 i# G. _<li>在我们这种情况下,self.workers只能作为引用被for循环迭代。</li>
, D0 u: H+ I1 [# N" Q, d</ol>. H+ N6 R! a, Y9 Z& l
<p>这里考虑让Worker持有<code>Option<JoinHandle<()>></code>,后续可以通过在Option上调用take方法将Some变体的值移出来,并在原来的位置留下None变体。<br>6 Q. T5 F. v6 R. j! s" ^
换而言之,让运行中的worker持有Some的变体,清理worker时,可以使用None替换掉Some,从而让Worker失去可以运行的线程</p>
; r# g" p6 a X8 G& ]! M<pre><code>struct Worker where. R! C7 F. ?$ a( M1 v9 l5 ?
{
2 R: a0 v$ `: k- j: d, G# z8 G) W _id: usize,
0 @+ E* M$ r9 w! O t: Option<JoinHandle<()>>,, _- Z' W8 a& _) t
}% d& P+ p4 j+ a, d: C
</code></pre>4 M H; y' B% G
<h1 id="要点总结">要点总结</h1>; O( K$ K( R, j6 s1 P: w0 O# f. I
<ul>/ `6 f0 O- B1 T6 z
<li>Mutex依赖于生命周期管理锁的释放,使用的时候需要注意是否逾期持有锁</li>
# V: o* E" @) ~8 n<li><code>Vec<Option<T>></code> 可以解决某些情况下需要T所有权的场景</li>
0 j, i' l3 `! b2 n. L6 D6 `/ q</ul>
x+ q% ~9 e9 N9 I) v3 N<h1 id="完整代码">完整代码</h1>
! k3 c! {5 M/ Z; h<pre><code>use std::thread::{self, JoinHandle};. C) j5 d1 d: T7 H, }
use std::sync::{Arc, mpsc, Mutex}; Y! b4 w4 X7 `- {9 i
( F U% Z j1 t
0 n. f0 V8 ?* f( ^" d( W2 Ptype Job = Box<dyn FnOnce() + 'static + Send>;+ T5 d) A C8 ^* Z2 o
enum Message {9 x/ r) l4 |: M% z6 h+ c+ Y
ByeBye,0 U" o9 v3 T2 F H4 p
NewJob(Job),
( x7 L/ G! z3 q% u: j6 F7 B}- f( E( [2 }- G
$ r3 Y, v0 `2 V* V- U! Y( E3 ^
struct Worker where
& H: @$ Z# ~# D% s# P$ G: N6 p{ R" k. K- L1 P z
_id: usize,. [. c# l* k; J5 }% U7 c' _
t: Option<JoinHandle<()>>,8 F' S9 w% t. a5 W$ K
}
/ A8 J2 R) m' p: m5 I
: _! u: E$ q Wimpl Worker! ~! w% u6 M3 u! O/ b. B
{
3 o* Y0 \/ z# X7 W& ^ X fn new(id: usize, receiver: Arc::<Mutex<mpsc::Receiver<Message>>>) -> Worker {
- v$ D7 J3 \2 V( x; P3 x/ j6 W let t = thread::spawn( move || {
6 j7 d; m, j4 ]. _1 z loop {
* L8 s! z0 N( |% _ let message = receiver.lock().unwrap().recv().unwrap();( o k9 h. G3 e) ^
match message {
7 ?' o8 L! w' J" l9 p Message::NewJob(job) => {
" ~( L& U" Y2 @2 c% s* f) o println!("do job from worker[{}]", id);5 O+ Y- S% w! y8 x8 ?" E6 x2 D3 m" `
job();* y0 I4 T* i1 [; B+ s1 U* u
},. e, z5 }) B. P0 `6 f
Message::ByeBye => {
7 |1 ?! \1 U! U6 v6 {7 g println!("ByeBye from worker[{}]", id);
3 R! T: v) B+ u4 c+ D" @0 ?7 g break9 [, T' q0 e) w+ p# x: L
},2 b& T5 X& E; C/ R
}
5 J* n1 H0 b; J: W9 B3 s }
9 F. E$ f/ f6 ~7 f+ N/ _ });' N0 v6 o5 f8 t2 a( d! s
$ K( \$ U3 T. `: e% M4 L0 T
Worker {
9 o; P7 }1 |0 ] _id: id,* Z |# `6 Y+ U2 p1 u
t: Some(t),/ i j% Z" B1 H/ X
}$ B0 Q0 p; u3 N; V, m# E
}$ e# @8 }5 x7 i$ E$ w8 u- K; T
}4 H5 b8 s' P, M; r! F5 V8 ?3 [
3 L" G+ f9 _1 Epub struct Pool {' E2 ^" g) F0 c. U m& Z; p
workers: Vec<Worker>,
/ {6 Q- U+ E9 O. N( o max_workers: usize,
& r1 ^5 y3 l4 u6 `; I3 o sender: mpsc::Sender<Message>0 h e7 g$ H7 ~- w
}
$ `# {& u# O! P/ |' Q: M8 M1 r0 F/ s
$ V0 m; A; v+ ]+ Ximpl Pool where {
% O, V. B1 \, r, P pub fn new(max_workers: usize) -> Pool {
1 v. l+ I* a% }* f x6 K if max_workers == 0 {
( Y: H A0 l4 n panic!("max_workers must be greater than zero!")4 @2 A# W8 X. _ m5 h/ i0 h& i
}* s# u8 l' Y! Z) D
let (tx, rx) = mpsc::channel();
4 w4 G4 {! u/ K) B8 b, x! I9 `& Q6 j6 k# h# C
let mut workers = Vec::with_capacity(max_workers);$ V2 \8 o" E, a. q5 w
let receiver = Arc::new(Mutex::new(rx));
7 ?0 P! ?7 @6 [( i' F% U for i in 0..max_workers {3 p3 I8 x# Q. v
workers.push(Worker::new(i, Arc::clone(&receiver)));
, H. ?+ i! _1 v$ G8 V }
3 \' n7 ~2 z' W7 Y& X& `" V
/ m. u( n1 l; q2 X* z Pool { workers: workers, max_workers: max_workers, sender: tx }$ [" {* ]' k" q4 i
}1 r7 F6 E+ a% s9 p
5 Y( B% f g0 N: s7 z y
pub fn execute<F>(&self, f:F) where F: FnOnce() + 'static + Send9 r0 e! U8 B. j# g% T( ?
{
6 K# d, b" G0 d: L' E
+ p Z0 b$ V" c$ q" ] let job = Message::NewJob(Box::new(f));
# i1 j, m& @0 l5 n self.sender.send(job).unwrap();( T! s( ~4 T. l* k& L( g3 ^0 |
}$ s! y6 m" c g- U$ ?6 X3 M7 j* W- v
}
7 i0 S; d% h, B* @4 }1 h
. {) c. \8 R# O. S7 w1 simpl Drop for Pool {
8 x' Z+ ]" h J+ L fn drop(&mut self) {
6 p/ \3 Q' F( U/ M; C3 g7 K for _ in 0..self.max_workers {' G. x) X" t) B5 |; \* E' Y- ?8 D
self.sender.send(Message::ByeBye).unwrap(); C/ G& G. B+ X* m2 C H
}' K$ s/ x; {% _1 y: }& I/ Z) \
for w in self.workers {
0 D6 N( g6 ~( t9 n m& L if let Some(t) = w.t.take() {1 i* p- H2 D: T3 r0 |
t.join().unwrap();. x: x: I1 H. I$ x
}( k9 H6 N/ L# c4 `
}
3 Y3 v/ M* o7 I7 v" S* \6 c: K }7 A/ n/ C8 U2 {
}
! \' G. O% e+ \ q" i. T
2 e$ K' [: u# w, J5 Y" v ?' J" h3 r& f
#[cfg(test)]
0 A8 }: \0 k9 `8 E! Qmod tests {* B! m& ]* p1 ^4 j8 k o3 H
use super::*;
- B' u. S2 k9 ^9 u- \% Y #[test]
3 Q8 M( D4 b3 `/ ?7 n: c$ M fn it_works() {
@3 @+ V" X1 {% i" a let p = Pool::new(4);3 x; y6 u7 q6 h$ a9 V
p.execute(|| println!("do new job1"));) W* k, K0 e/ K1 j% j
p.execute(|| println!("do new job2"));
. G# B6 c8 b& K6 W p.execute(|| println!("do new job3"));/ j( I2 p& x! \ h$ D; x
p.execute(|| println!("do new job4"));3 ^: K# l3 Y3 O
}. Y% e5 T$ S! Q7 l9 J+ `$ C
}4 l3 z4 t6 |7 l2 o4 [5 ]7 r
</code></pre>- j4 i7 p( i' l# O
0 o4 R. v0 C# Z* z# H. U |
|