|
|
# V9 A# B8 P% y# P7 _<h1 id="如何实现一个线程池">如何实现一个线程池</h1>
2 c: _1 d4 W. E: o+ g# n<p>线程池:一种线程使用模式。线程过多会带来调度开销,进而影响缓存局部性和整体性能。而线程池维护着多个线程,等待着监督管理者分配可并发执行的任务。这避免了在处理短时间任务时创建与销毁线程的代价。线程池不仅能够保证内核的充分利用,还能防止过分调度。可用线程数量应该取决于可用的并发处理器、处理器内核、内存、网络sockets等的数量。 例如,对于计算密集型任务,线程数一般取cpu数量+2比较合适,线程数过多会导致额外的线程切换开销。</p>
* v5 W6 b& L$ G2 ]1 k& Q<p>如何定义线程池Pool呢,首先最大线程数量肯定要作为线程池的一个属性,并且在new Pool时创建指定的线程。</p>
5 x9 J g" l8 i5 L' I" Z<p>线程池Pool</p>
8 T. e, n9 I2 G# a4 C<pre><code>pub struct Pool {- j; e, A& o1 J* G$ c4 N# b1 R: G
max_workers: usize, // 定义最大线程数0 m- x" c" [' B( ?9 z
}. \( p* F( s2 T) \7 {; r0 F8 C5 |
2 z* `# s+ j* {5 ~$ P( Simpl Pool {7 z" h/ [0 M" U1 y9 `
fn new(max_workers: usize) -> Pool {}
+ ^' G3 r4 ]) g+ _4 \5 \) X fn execute<F>(&self, f:F) where F: FnOnce() + 'static + Send {}
! b9 g+ k6 z+ B1 z4 j2 u2 x2 p# J' A}6 [/ L# ?7 L W2 Y( c J8 q
3 B4 e: V# l) O2 q! v+ _
</code></pre>& k; R' K& Q$ R. s; p5 P2 Q
<p>用<code>execute</code>来执行任务,<code>F: FnOnce() + 'static + Send</code> 是使用thread::spawn线程执行需要满足的trait, 代表F是一个能在线程里执行的闭包函数。</p>
. E% ^8 x5 `' @0 H2 R/ u _8 f<p>另一点自然而然会想到在Pool添加一个线程数组, 这个线程数组就是用来执行任务的。比如<code>Vec<Thread></code> balabala。这里的线程是活的,是一个个不断接受任务然后执行的实体。<br>1 \: g; { z5 J, Q* C* C& ~
可以看作在一个线程里不断执行获取任务并执行的Worker。</p>0 i% z# G8 m: }% f: \, w( x
<pre><code>struct Worker where6 j* E t- U$ j: o% J s/ S' b
{
# `# A' i9 p( I% R+ g7 G3 t4 m _id: usize, // worker 编号- u4 m# x7 P! C
}
: }/ q, r% M" B. i</code></pre>4 U5 ?: G8 ]8 p' E0 b, e4 e
<p>要怎么把任务发送给Worker执行呢?mpsc(multi producer single consumer) 多生产者单消费者可以满足我们的需求,<code>let (tx, rx) = mpsc::channel()</code> 可以获取到一对发送端和接收端。<br>9 t3 G# K5 W9 j5 ~! @7 h2 {- V, N1 d
把发送端添加到Pool里面,把接收端添加到Worker里面。Pool通过channel将任务发送给多个worker消费执行。</p>7 w$ B8 T, q/ V0 ^
<p><strong>这里有一点需要特别注意,channel的接收端receiver需要安全的在多个线程间共享</strong>,因此需要用<code>Arc<Mutex::<T>></code>来包裹起来,也就是用锁来解决并发冲突。</p>
' T5 \0 z% E6 U* s, w<p>Pool的完整定义</p>
/ r' `' N5 t! i. p% w; I<pre><code>pub struct Pool {
3 t$ v6 {2 K. T5 ?+ m workers: Vec<Worker>,4 q* J9 B F8 \
max_workers: usize,
7 A: [/ H5 E2 u9 h% L# s; s+ [ sender: mpsc::Sender<Message>7 u) X/ d! c4 y4 h. N% l
}' @* d/ X* e7 H3 H& K5 r
</code></pre>1 }+ [- R0 s! s! G7 M% |5 A* z! S
<p>该是时候定义我们要发给Worker的消息Message了<br>
: [6 b" v1 b5 y3 {7 t定义如下的枚举值</p>
; O k" B5 `* T8 s- ]4 c<pre><code>type Job = Box<dyn FnOnce() + 'static + Send>;
$ |' y N. o& m9 d3 Yenum Message {* H, |8 v8 h9 j6 P0 P
ByeBye,& l6 N9 u/ d, S& Q" p4 O
NewJob(Job),4 |" Y# I. l. m
}+ y- o2 v! x a+ ?) m
</code></pre>/ u+ q; \* Q2 r( o4 ] h/ E
<p>Job是一个要发送给Worker执行的闭包函数,这里ByeBye用来通知Worker可以终止当前的执行,退出线程。</p>$ \; q# M2 `! C, Y1 _2 C S* |0 ^
<p>只剩下实现Worker和Pool的具体逻辑了。</p>1 B/ M2 x2 j7 F2 E8 V7 e
<p>Worker的实现</p>; `* W6 F% B8 P4 L* b/ S3 {
<pre><code>impl Worker9 h3 c+ Z4 d, P# S( f6 t# `
{; [/ K' t4 b* i
fn new(id: usize, receiver: Arc::<Mutex<mpsc::Receiver<Message>>>) -> Worker {* u5 p1 Q- J% b
let t = thread::spawn( move || {
# d6 v# n; n M5 C loop {
9 t# G+ f; B" D" n/ h let receiver = receiver.lock().unwrap();2 h" v# K5 O. d6 ]- m
let message= receiver.recv().unwrap();
$ X; R* L4 w: ^" Q& U match message {
1 B/ b- V- w H9 J4 C/ E7 ]6 o0 D- r Message::NewJob(job) => {1 [% f# R) B o5 H
println!("do job from worker[{}]", id);
+ {, |- j8 d+ Y( }4 I job();
y8 W! I% E' X0 @7 W },
6 P& \4 L" g [' s5 u q Message::ByeBye => {6 f! g. F9 i/ m+ h* @
println!("ByeBye from worker[{}]", id);
4 W7 p A# Y0 W9 H% _0 c break
8 D/ W0 i7 t: d D; Q5 d* }- B& U, s },3 {3 _; A* b& H5 x# Q( R$ ?
} + C# O$ p9 J1 z( ?& i+ f
}
. D5 X$ Z& t! T0 Q B' Z& j });
$ `' s: g- E# w, R7 K% r) w* b2 |/ T0 F8 b
Worker {: t5 \( c$ G& W0 r
_id: id,6 m+ ]# x0 p8 O" d4 o+ t1 e
t: Some(t),% z9 h1 [3 o7 U7 J$ I5 t
}: q Y6 [7 e3 ] R
}
: t( h. V; t0 h' k7 Z& v! E}
3 u. Y( O) C2 [8 h" C9 @8 Z</code></pre>+ n, }" F! {0 H! A3 x: R6 u7 n
<p><strong>let message = receiver.lock().unwrap().recv().unwrap();</strong> 这里获取锁后从receiver获取到消息体,然后let message结束后rust的生命周期会自动释放掉锁。<br>
% T1 f3 \/ `* y7 `但如果写成</p>
; j3 z% z9 T! Z2 t' _% m<pre><code>while let message = receiver.lock().unwrap().recv().unwrap() {1 ?, a$ o4 H3 j n, d
};
1 }; V1 Y2 m/ o. w</code></pre># j9 l8 k- _0 q$ q4 | A; y" Z) L
<p>while let 后面整个括号都是一个作用域,要在这个作用域结束后,锁才会释放,比上面let message要锁定久时间。<br> m) L! o. l! g, W+ |1 Y7 |) f& R
rust的mutex锁没有对应的unlock方法,由mutex的生命周期管理。</p>
( J& c' g- x* {( I* I<p>我们给Pool实现<code>Drop</code> trait, 让Pool被销毁时,自动暂停掉worker线程的执行。</p>3 O+ r( Q: M6 V9 I$ F0 `
<pre><code>impl Drop for Pool {6 k/ v- Y) Z! h+ [. o! K$ F; v
fn drop(&mut self) {3 o( Y3 e$ k, U9 S) ]
for _ in 0..self.max_workers {
& R$ e% T( t' Z self.sender.send(Message::ByeBye).unwrap();
) Q+ ^0 Z/ Q4 S+ h }7 U/ `( \, y+ X9 `
for w in self.workers.iter_mut() {+ _4 G9 [+ S5 o- @1 V/ _4 |- ^
if let Some(t) = w.t.take() {0 D+ S' J7 [# L$ X. C
t.join().unwrap();( f: v f- r) N2 |6 N
}
+ w4 _: F- ?6 k% | }' \, L% e4 I: w3 g8 u) ^+ p" W
}' o, {9 n+ I6 Q/ u3 P
}
4 K6 w2 _: |# T$ [- s2 L- M' k$ [* B7 n! c! m7 P& M
</code></pre>
# l, A/ H$ s$ a6 V' [<p><strong>drop方法里面用了两个循环</strong>,而不是在一个循环里做完两件事?</p>5 {# X$ d' N) b7 [; I- i. k& V
<pre><code>for w in self.workers.iter_mut() {
1 p( D$ W8 a. v" d; H if let Some(t) = w.t.take() {
+ m# l% q+ W' i$ [ self.sender.send(Message::ByeBye).unwrap();
4 {2 Q8 z$ Z7 C( i t.join().unwrap();
. R/ u/ D. N5 V' A7 _/ v }2 _* o3 d5 D6 h, n, O7 _
}, c( G' u; V. [2 O
W5 Z" t+ [. n/ f
</code></pre>
/ A" O/ Z+ T" }<p>这里面隐藏了一个会造成死锁的陷阱,比如两个Worker, 在单个循环里面迭代所有Worker,再将终止信息发送给通道后,直接调用join,<br>
3 s, n) F, }5 y2 e5 H我们预期是第一个worker要收到消息,并且等他执行完。当情况可能是第二个worker获取到了消息,第一个worker没有获取到,那接下来的join就会阻塞造成死锁。</p>
" G1 L: m+ F6 Z. a<p><strong>注意到没有,Worker是被包装在Option内的</strong>,这里有两个点需要注意</p>
: k0 z; S4 ]) l<ol>/ e7 W0 k2 f% e6 a
<li>t.join 需要持有t的所有权</li>
$ }( ]5 z0 Z" e; @0 Z, g<li>在我们这种情况下,self.workers只能作为引用被for循环迭代。</li>
: l, `5 c3 W- p) s7 r</ol>
7 B* e& z3 K) H<p>这里考虑让Worker持有<code>Option<JoinHandle<()>></code>,后续可以通过在Option上调用take方法将Some变体的值移出来,并在原来的位置留下None变体。<br>
6 `5 U; k9 u( z9 R+ V( w/ D( ?换而言之,让运行中的worker持有Some的变体,清理worker时,可以使用None替换掉Some,从而让Worker失去可以运行的线程</p>$ @* N. x( R* b- I' |6 n
<pre><code>struct Worker where5 w' O% A" A) M/ o8 b( k
{
( N) }/ ^0 f# ~ _id: usize,
: w7 N* k1 H& Q t: Option<JoinHandle<()>>,
1 S) k/ h! X+ T- i9 X}, _* |: i8 z; D [
</code></pre>% ^) ~" S' K& x5 e
<h1 id="要点总结">要点总结</h1>4 `) r, n3 x, w& X
<ul>
$ d" Z# c, h+ D& z) \& k0 W<li>Mutex依赖于生命周期管理锁的释放,使用的时候需要注意是否逾期持有锁</li>/ U! t) K* l; b8 V/ x$ u& J2 v
<li><code>Vec<Option<T>></code> 可以解决某些情况下需要T所有权的场景</li>" b: L8 Q; d" S" |: _9 n6 M; S
</ul>, m% E9 q" }+ q* I2 Q0 g, D0 c
<h1 id="完整代码">完整代码</h1>
* V( Z2 N" |& D! m8 N<pre><code>use std::thread::{self, JoinHandle};
. E) c4 Y& D5 |( J; M [! j7 nuse std::sync::{Arc, mpsc, Mutex};
8 k; k" S. M3 }* p; X3 w
3 c9 j0 e1 f* u/ L/ y# v
8 n) H; j7 ?/ x0 `type Job = Box<dyn FnOnce() + 'static + Send>;
2 y7 ?0 ? c4 @3 i! v2 n: w9 N4 Menum Message {* t r+ E: T; V6 I1 ]% n
ByeBye,, K; ?' E$ Z9 r: Y9 ?0 @$ f
NewJob(Job),
; W+ ~$ V$ ~0 [! D3 N}
8 g, X% E* Y3 `% {( W/ A7 j) A; c7 c' z# l
struct Worker where
) T5 e# C+ I% a% Y: \{; v* f7 s% d/ ?- T; T. L- [
_id: usize,
# f( j4 t. c, e t: Option<JoinHandle<()>>,
' n4 @2 p3 C! M$ C: ^}7 g. z* ?, ]" r6 P8 A6 y9 j
( N8 A0 R& U" p; E" w$ [. A# t6 W: z
impl Worker
; e( w2 q; ^5 G! d# {2 E{' x# i2 P- ^6 z% L9 @& }
fn new(id: usize, receiver: Arc::<Mutex<mpsc::Receiver<Message>>>) -> Worker {# ~: Y- `8 ^" Z% A
let t = thread::spawn( move || { R1 Z/ m& g6 | b% e
loop {
. d% p1 h* k1 l0 J let message = receiver.lock().unwrap().recv().unwrap();
p0 O6 G) ` t% H match message {7 E- V. }. I1 A* Y5 L) m
Message::NewJob(job) => {
- E6 B5 t! X4 q7 u V println!("do job from worker[{}]", id);0 {! P) B7 q/ F% y9 J. @" {
job();
; `7 F& E- k+ C },
+ T6 f4 V4 M; O5 S2 W' Y( X Message::ByeBye => {
0 Q U* ~3 Z7 G* x: k+ g println!("ByeBye from worker[{}]", id);5 w+ q! K! \# N9 ^
break
8 t! K( w' j( _ }, G4 H) y* F' m4 U0 _
} . x, W7 d6 _' L: w8 M, |; b
}
: k6 Y# y1 _( Q5 e });) f @$ N* I; |9 F) D1 |% U) G% A
6 N+ V! X/ C; v% O1 U Worker {0 p% R7 N; O- G. {: ]; N" L* Y8 b
_id: id,1 @) r" q S6 T. w
t: Some(t),
$ H0 g" M& w8 j: G/ M4 B8 v0 | }* n$ t& g, d( M: F
}$ s" E- \( ]+ }) D
}
6 W" o9 L% @, J: t* l' D+ \, O, k
% n1 E" a: w3 E; spub struct Pool {
g ~0 r0 q0 L+ }8 e workers: Vec<Worker>,
2 U/ y" n3 _6 F9 O# y( N9 j$ q0 v2 ^0 q max_workers: usize,
. ^6 y4 c ~! `+ h sender: mpsc::Sender<Message>
! Y; v- S# k& c. h e4 R}9 v" V) [7 q- Z1 j& s4 `
5 ^# L: F! T6 S0 o9 E W1 _impl Pool where {* a ? P8 V9 c2 W
pub fn new(max_workers: usize) -> Pool {8 ^# n) z+ `) u$ ~6 D
if max_workers == 0 {
2 h+ a. a9 l( e3 ~ panic!("max_workers must be greater than zero!")2 t7 i/ j% A; G1 m. R. N- I9 _
}- o; d0 m( e1 n2 a5 c
let (tx, rx) = mpsc::channel();
, G; A! Z; y- u7 v" P% {4 ]3 A' {) \: G1 M/ n
let mut workers = Vec::with_capacity(max_workers);
7 N' R- R- [/ [" b: w% T let receiver = Arc::new(Mutex::new(rx));
0 d; t# G: H9 s8 O4 N for i in 0..max_workers {
[+ w. G# k7 |7 L8 \; @/ b6 } workers.push(Worker::new(i, Arc::clone(&receiver)));" ?8 H& F( o( F# h
}
" O. J" _) E3 x( c
# L$ Q: p- ~2 \" v. o Pool { workers: workers, max_workers: max_workers, sender: tx } [; {0 T! p' q
}
" I+ C( g& N0 U ; {8 p: ]1 H7 N/ e# y
pub fn execute<F>(&self, f:F) where F: FnOnce() + 'static + Send z9 h( [% T0 x1 I! V! B
{1 r- Z# X; B% N
; C& I9 x# }" M. p let job = Message::NewJob(Box::new(f));
: g% r3 m* n$ p, f3 D self.sender.send(job).unwrap();* e& c8 T7 _% J6 B
}0 G/ {) ~* E7 f: A; g' k$ F$ s/ t- @
}/ R! S- j) g8 ]% D
4 o0 E" h& l2 C! t4 i
impl Drop for Pool {
$ ]$ [# l1 y5 e, Y1 J fn drop(&mut self) {
1 h- e! ^ U3 _, x for _ in 0..self.max_workers {
2 m3 {8 ~8 P2 U& R self.sender.send(Message::ByeBye).unwrap();3 J. s, U- D* M
}
( x4 `& i. B! m/ y% }- ?8 B+ D) C5 c, Q for w in self.workers { I& `4 u# T; N& E
if let Some(t) = w.t.take() {1 @" z$ _* q+ {0 Z/ B
t.join().unwrap();/ L, m' M+ M# z7 h2 i) _
}
3 w" O! Q* y' k6 h/ X/ P: H }# D6 _6 u8 p2 C z
}& `' v3 n" \2 p/ @! g
}
; O- L" j: d# i5 T* O3 d
# N. x0 |) X% z5 q$ i8 t9 s% R* t* a$ x* e" L" r# {
#[cfg(test)]
% v( u9 V, F4 u3 f$ ~mod tests {( r, i: T: s9 g
use super::*;' k& j, @" S6 C' O' e5 @
#[test]* [5 {+ ^7 R- b
fn it_works() {% R, j- R" Y, c$ O
let p = Pool::new(4);: a5 e C* a0 N) v/ P& d, t J6 w. r
p.execute(|| println!("do new job1"));2 N) X$ S4 J4 n% {8 X7 V
p.execute(|| println!("do new job2"));
& G- ]% a1 N3 ]& b |; ~6 E p.execute(|| println!("do new job3"));
6 t% J% O3 E: v% \* U. t9 R p.execute(|| println!("do new job4"));
; m$ a; w, u& Y$ p/ t6 ^; C8 e }, G" j% ^, P% f! S O7 R! |$ E
}8 t' |3 d, K; o0 j) c6 [& @
</code></pre>
4 R& X% Y, V/ s4 \. j j+ F. ?6 Z
, Z1 b6 W n1 Y |
|