飞雪团队

 找回密码
 立即注册
搜索
热搜: 活动 交友 discuz
查看: 19416|回复: 0

rust 实战 - 实现一个线程工作池 ThreadPool

[复制链接]

9171

主题

9259

帖子

2万

积分

管理员

Rank: 9Rank: 9Rank: 9

积分
29843
发表于 2022-2-12 14:35:42 | 显示全部楼层 |阅读模式
1 b- e7 G5 B7 |& z4 m
<h1 id="如何实现一个线程池">如何实现一个线程池</h1>
+ y: m% j0 v* X" z! d2 S<p>线程池:一种线程使用模式。线程过多会带来调度开销,进而影响缓存局部性和整体性能。而线程池维护着多个线程,等待着监督管理者分配可并发执行的任务。这避免了在处理短时间任务时创建与销毁线程的代价。线程池不仅能够保证内核的充分利用,还能防止过分调度。可用线程数量应该取决于可用的并发处理器、处理器内核、内存、网络sockets等的数量。 例如,对于计算密集型任务,线程数一般取cpu数量+2比较合适,线程数过多会导致额外的线程切换开销。</p>
+ M5 K' W3 U! L3 Q4 Y" o" T<p>如何定义线程池Pool呢,首先最大线程数量肯定要作为线程池的一个属性,并且在new Pool时创建指定的线程。</p>" S; u& V* `6 u3 m! Q
<p>线程池Pool</p>) U! }: ~! ^' F/ X. }
<pre><code>pub struct Pool {$ j/ W! i% `& `- z6 t) w
  max_workers: usize, // 定义最大线程数, w7 E% ^& ^$ T$ r; l
}. I" W: J: D% S) D( n4 _4 q, T+ Z
" n) z0 H5 `7 _* {( c: H" T  n
impl Pool {3 j' o9 o+ U3 c
  fn new(max_workers: usize) -&gt; Pool {}
! ?. ^& K0 E& V' A+ u) {) P  fn execute&lt;F&gt;(&amp;self, f:F) where F: FnOnce() + 'static + Send {}
! u" [( z5 ]2 F& H% \# f0 u/ i}
/ E! @  j( X9 K8 b: g, \( H9 @  i) k( i" {
, v6 O: W7 q7 m2 O4 Y5 g4 X- P' X% m</code></pre>
7 b' \) a' S6 \! X<p>用<code>execute</code>来执行任务,<code>F: FnOnce() + 'static + Send</code> 是使用thread::spawn线程执行需要满足的trait, 代表F是一个能在线程里执行的闭包函数。</p>
2 I! p) r4 N) u9 y) b! U% n<p>另一点自然而然会想到在Pool添加一个线程数组, 这个线程数组就是用来执行任务的。比如<code>Vec&lt;Thread&gt;</code> balabala。这里的线程是活的,是一个个不断接受任务然后执行的实体。<br>5 W3 r+ r/ ~# H/ a
可以看作在一个线程里不断执行获取任务并执行的Worker。</p>
! g6 w( K/ T6 [# v; ^+ ?* _<pre><code>struct Worker where
( K$ O1 o* z0 L: S1 J{
8 Y) C1 q% I0 G    _id: usize, // worker 编号/ F: _. r5 V% D
}
& ^( Q- O; O4 _</code></pre>
& `" ~5 P4 o7 l4 H$ {<p>要怎么把任务发送给Worker执行呢?mpsc(multi producer single consumer) 多生产者单消费者可以满足我们的需求,<code>let (tx, rx) = mpsc::channel()</code> 可以获取到一对发送端和接收端。<br>
. }, m3 s! D) u8 k# G3 I6 K把发送端添加到Pool里面,把接收端添加到Worker里面。Pool通过channel将任务发送给多个worker消费执行。</p>
, G4 y2 j% C5 K5 N. Y<p><strong>这里有一点需要特别注意,channel的接收端receiver需要安全的在多个线程间共享</strong>,因此需要用<code>Arc&lt;Mutex::&lt;T&gt;&gt;</code>来包裹起来,也就是用锁来解决并发冲突。</p>
3 a1 m7 u+ @6 s$ ]1 ]" p3 R<p>Pool的完整定义</p>' @; i! U3 Q4 l. A4 k4 P7 z
<pre><code>pub struct Pool {: @, I) }* ?8 i7 Y6 d, m1 w
    workers: Vec&lt;Worker&gt;,
9 e; P; C6 L" l/ G/ m    max_workers: usize,1 Y& x2 {: W: G
    sender: mpsc::Sender&lt;Message&gt;- @' i1 m& X" {+ N$ C; l
}
# Q) j) Z: [' {</code></pre>
: m, H! L! T$ ~' C, w( D<p>该是时候定义我们要发给Worker的消息Message了<br>: R' p& ^+ a( _
定义如下的枚举值</p>
5 ]" B# B& k; M" D" Q4 h6 y) j, q<pre><code>type Job = Box&lt;dyn FnOnce() + 'static + Send&gt;;( C& r: N9 P, z8 D" {- h
enum Message {
5 e4 F) {( q$ H+ G; f1 Z+ a, H, }  j# G    ByeBye,( m8 k3 T8 `# W; G) G) N
    NewJob(Job),: a" O6 [8 m: c" J
}
% |, I+ w7 u5 u* F</code></pre>) y* W( @0 o0 O2 A, l4 v. I
<p>Job是一个要发送给Worker执行的闭包函数,这里ByeBye用来通知Worker可以终止当前的执行,退出线程。</p>
- d3 K8 x* @/ `$ }<p>只剩下实现Worker和Pool的具体逻辑了。</p>+ g8 e/ }$ `+ q& e- I* G; {
<p>Worker的实现</p>
" O3 |2 D2 K' Y8 w1 P6 Q<pre><code>impl Worker
4 t" `) v" k( I; S{
3 t  o5 i9 O! i9 i    fn new(id: usize, receiver: Arc::&lt;Mutex&lt;mpsc::Receiver&lt;Message&gt;&gt;&gt;) -&gt; Worker {
9 k  {, }, A) ]& k" x. J        let t = thread::spawn( move || {% P9 O  E7 a5 A
            loop {, \4 V: \: [" N5 Z% j% ~* K- M
                let receiver = receiver.lock().unwrap();
. y5 o+ b; M! i! S, ^. ]2 O( J8 W                let message=  receiver.recv().unwrap();
4 p: A& ?) J! J2 E, b. X                match message {
. ~" M4 J# b) O. K. e3 r                    Message::NewJob(job) =&gt; {) U* i2 x, [/ I
                        println!("do job from worker[{}]", id);
* ]" F; k3 P# @. K+ ?" `  e                        job();: q9 H2 S5 h0 j) E: v$ C
                    },
0 m1 I3 ]  ]5 R4 i                    Message::ByeBye =&gt; {
5 b1 b+ }) ~2 B+ u4 r# G- u: q                        println!("ByeBye from worker[{}]", id);
! f+ x% w/ {, ~: T                        break7 C1 }& [, p/ H( A$ E) ]! i
                    },
/ u2 t5 u" `9 O- L: w$ l' d                }  
+ v& O! z& o/ i& k* t            }
$ d: ]0 {# k9 G; m1 ^& I. x        });
  a! ]# ], x" I* |* ?; N8 P0 }. c- ?/ ]6 S0 r
        Worker {2 k% R3 Q$ L9 H3 C1 ]: A5 ^' {
            _id: id,; p# B6 _: F( W2 x3 \2 x, K; o
            t: Some(t),
7 G) u6 A5 r  c* R2 R        }* \- m$ }* u9 R4 n7 R; ]2 ^' F( H
    }- `6 E0 V- d+ ^: x; m: A
}
! r7 U- G7 b$ y8 w( l4 g6 W, f</code></pre>
# R) {" G: b  \/ c& ?. O/ v. T<p><strong>let message = receiver.lock().unwrap().recv().unwrap();</strong> 这里获取锁后从receiver获取到消息体,然后let message结束后rust的生命周期会自动释放掉锁。<br>& m+ _6 X- K+ P1 E
但如果写成</p>' R. ]0 z! I: B  Y/ J% x' w
<pre><code>while let message = receiver.lock().unwrap().recv().unwrap() {
6 y/ N. C1 O) ]( l8 n- d};: ^5 [( ^" q6 T% e
</code></pre>. j/ p3 K/ T* y9 X6 l2 B
<p>while let 后面整个括号都是一个作用域,要在这个作用域结束后,锁才会释放,比上面let message要锁定久时间。<br>+ _" _3 C) ?7 \3 j3 p1 \
rust的mutex锁没有对应的unlock方法,由mutex的生命周期管理。</p>0 H  w1 u; v; P% l' }
<p>我们给Pool实现<code>Drop</code> trait, 让Pool被销毁时,自动暂停掉worker线程的执行。</p>
" `+ I) ^' p7 W0 Z+ w<pre><code>impl Drop for Pool {$ T2 [- X5 c' E* p
    fn drop(&amp;mut self) {0 N# L- D0 i0 C# y# H
        for _ in 0..self.max_workers {
% a4 X3 [0 z% {. Z, d            self.sender.send(Message::ByeBye).unwrap();+ X: F* ~5 H0 L0 I! J* [7 J
        }
0 a6 n: F  A2 z6 a7 G* ^' ^- g        for w in self.workers.iter_mut() {
( l2 G' B9 ^" \2 W* `            if let Some(t) = w.t.take() {
5 y9 ^: |  Q: I$ N- \/ O; c                t.join().unwrap();1 f/ k) K3 D' T4 \
            }( [& F! E( D. Z- `+ S
        }
! p( k% |4 R$ b; x5 h    }
; ^" U( y9 M8 R+ [4 J4 C}
# |! s$ m) t4 [+ t8 g6 m* O0 Y/ a0 j$ D
</code></pre>
+ U$ x: D. S: S; H5 r+ `8 i<p><strong>drop方法里面用了两个循环</strong>,而不是在一个循环里做完两件事?</p>" s/ N+ }2 q4 B" R( o- Z
<pre><code>for w in self.workers.iter_mut() {( M( M1 s; F1 T, y; O
    if let Some(t) = w.t.take() {
$ L8 Y  W) j4 q2 z) L2 p5 s/ a        self.sender.send(Message::ByeBye).unwrap();
  l& i- W& E% c5 Y( S" K5 _9 f        t.join().unwrap();
7 Y, S. h) H2 T% O    }
; n# W" y7 m  A& X4 z; F9 T* _  r3 ~}
! Q9 m3 U4 K' A1 H4 m4 H$ i, {" w9 a* h) }) Y: R4 M
</code></pre>
8 |  }# j! O& ^; ?& {. O<p>这里面隐藏了一个会造成死锁的陷阱,比如两个Worker, 在单个循环里面迭代所有Worker,再将终止信息发送给通道后,直接调用join,<br>
4 j4 U  A" h7 B. D/ }3 D) }我们预期是第一个worker要收到消息,并且等他执行完。当情况可能是第二个worker获取到了消息,第一个worker没有获取到,那接下来的join就会阻塞造成死锁。</p>  X- g& `) @9 R( t) o/ L
<p><strong>注意到没有,Worker是被包装在Option内的</strong>,这里有两个点需要注意</p>0 h4 N' H7 i9 Z  I
<ol>
7 z1 j  ], u  @. {1 E* W<li>t.join 需要持有t的所有权</li>5 l9 V( f6 l( B% S9 ^! _/ H. R* V; U
<li>在我们这种情况下,self.workers只能作为引用被for循环迭代。</li>1 y' y6 n: j( m/ V0 a$ S
</ol>
! x% a! j: o& G' }# e<p>这里考虑让Worker持有<code>Option&lt;JoinHandle&lt;()&gt;&gt;</code>,后续可以通过在Option上调用take方法将Some变体的值移出来,并在原来的位置留下None变体。<br>
9 w1 G0 G  i/ V- P. S! `  L6 @8 e换而言之,让运行中的worker持有Some的变体,清理worker时,可以使用None替换掉Some,从而让Worker失去可以运行的线程</p>% i$ X- w+ b3 q- M* Y( o
<pre><code>struct Worker where- R5 B6 e! ^/ I  ~/ W- O& i
{% c1 ^+ {3 C. r3 F; [$ l# P
    _id: usize,, t% k* Q$ K- h& k. N3 c
    t: Option&lt;JoinHandle&lt;()&gt;&gt;,' U1 ~( f' M, X; h# J" r
}
- u$ {1 `# b7 M* d7 x</code></pre>0 T0 I2 T7 A9 h# N  p" w# y
<h1 id="要点总结">要点总结</h1>
- \# X5 U! r' S) Q* D% k" @# W4 k) l<ul>( {( r3 Z5 ^+ U/ x7 C
<li>Mutex依赖于生命周期管理锁的释放,使用的时候需要注意是否逾期持有锁</li>4 Z/ V; v) {! c/ Y2 |( }- M
<li><code>Vec&lt;Option&lt;T&gt;&gt;</code> 可以解决某些情况下需要T所有权的场景</li>2 N, y; D7 j% [3 O) q/ J; S, r( ]+ j& q
</ul>
. [' a  }: I' H0 N$ c4 ^<h1 id="完整代码">完整代码</h1>4 c0 O2 N' O: K5 V- Z- a
<pre><code>use std::thread::{self, JoinHandle};
4 {4 M1 ]) R9 e* h* @/ a. ]use std::sync::{Arc, mpsc, Mutex};
8 \4 U! z/ `* \' v5 ^* {8 a+ W0 N  g6 a

0 Z! ?# y. i3 ztype Job = Box&lt;dyn FnOnce() + 'static + Send&gt;;
# ?* {' v* H, V' U* e- C6 |& @# e3 }enum Message {& x9 f. Y) F. R
    ByeBye,. r; H( N7 x2 \; F, r2 i
    NewJob(Job),+ |( E6 G6 e1 _/ P
}) T: z7 B# p" n  E

6 ]& H  A3 G* Lstruct Worker where
% t9 _* M0 H, G  z{3 X3 ]6 `: H. V. S; g
    _id: usize,
( M8 p$ X) z0 W! |4 z    t: Option&lt;JoinHandle&lt;()&gt;&gt;,
2 G, d6 C4 K$ a- U! i5 k7 x; w}
* Z5 K' U: `& I. C  ~7 x' ^6 ~( {7 |) Q
impl Worker
+ ~$ b$ N: G7 K7 V# H{
$ g1 Y  ]8 {" W1 g  \    fn new(id: usize, receiver: Arc::&lt;Mutex&lt;mpsc::Receiver&lt;Message&gt;&gt;&gt;) -&gt; Worker {) F: v2 Z$ s  s6 a6 `* k# k' r: ~1 W5 D
        let t = thread::spawn( move || {! y/ G. D# w) s( `8 @
            loop {
" Z0 ?8 \% K: V; c4 v$ d. G                let message = receiver.lock().unwrap().recv().unwrap();
3 e+ t% z8 P) }) z. {2 X                match message {# x5 w/ \" A) D7 k
                    Message::NewJob(job) =&gt; {( _* ?# E4 I5 _" J) Q- ^! |* K5 Q
                        println!("do job from worker[{}]", id);
) j( k+ w* ]) W" }' W  v                        job();
% ~! ?$ f& \  m6 J3 g  r& u1 t                    },+ o# o2 ~( G+ {  n- x
                    Message::ByeBye =&gt; {' b; Y0 Q5 X) L  G( ~
                        println!("ByeBye from worker[{}]", id);
9 M+ A- ~4 t( T( e                        break9 h8 B0 L/ [* y! W
                    },2 _! p% C3 i9 J2 f  @" b
                }  
0 o  R0 u5 S3 |9 I1 W$ ~0 D            }& ?5 h* n) c; [
        });. w$ y1 P( v; a

. Y8 f2 m8 j6 X# Y( a6 o7 _' d        Worker {
4 ^$ N( P  @7 d* h2 T            _id: id,/ n& d6 x0 m  B" ?1 ^3 w' b, Z
            t: Some(t),: l5 M4 N( a: e' q
        }4 d" {7 a! `( k* {( }
    }
, b% m9 ~8 N) u/ B}8 s3 j- i8 d' S7 a3 A5 l, u- w- {

0 a% j9 H7 O- e) z  A' Npub struct Pool {- }5 |$ I  q, W
    workers: Vec&lt;Worker&gt;,
! r$ E* ~; Z3 {: B; a" j3 I% O    max_workers: usize,
8 |9 ^  L" k# U" o; ?" A( d8 _    sender: mpsc::Sender&lt;Message&gt;
) a7 j+ u. `* |; m}
$ N+ s  O0 ?, j
/ h$ T1 J% b- E3 @impl Pool where {
- T( ^2 P1 V, g    pub fn new(max_workers: usize) -&gt; Pool {- t; U6 y) d' Y+ |0 I7 `
        if max_workers == 0 {
) E- V8 W5 p& O! f: g            panic!("max_workers must be greater than zero!")
1 Z4 E' E+ f. y: H2 R. b8 ?        }
. U0 c5 T) |8 L* j8 x+ \        let (tx, rx) = mpsc::channel();7 j2 d7 H4 K0 }6 u% ~+ r1 t

3 S6 }, Z- g8 e8 K: a9 Y( h        let mut workers = Vec::with_capacity(max_workers);  ~$ S! S. N5 m9 g
        let receiver = Arc::new(Mutex::new(rx));
3 J' o* m# W' J0 E6 G        for i in 0..max_workers {2 N5 l' ?5 X- a' G! d
            workers.push(Worker::new(i, Arc::clone(&amp;receiver)));$ J# U1 k3 G6 O7 Z: S$ s
        }
. d) L7 v! |8 i3 t/ r$ J3 @4 g/ N3 V7 ]3 j, C# l- P$ o
        Pool { workers: workers, max_workers: max_workers, sender: tx }4 X$ p- o* e7 Z
    }2 o5 u' Z1 H$ O0 R
   
3 s8 F' i  I5 l    pub fn execute&lt;F&gt;(&amp;self, f:F) where F: FnOnce() + 'static + Send2 D! e+ ?% z6 c; Z$ h6 s
    {
1 T* ^( K! I- Q5 [! j% a6 M5 K; v8 `  |% K* x: @
        let job = Message::NewJob(Box::new(f));: e' G6 k- x: O8 o' {+ P
        self.sender.send(job).unwrap();
1 B+ L% e! t, U( ^    }, [# y; l0 U- P  `5 q
}
1 p" L6 \% p) `8 n! j: h9 K: _
0 s/ A, C! u# k% j2 ximpl Drop for Pool {
+ M3 I# I5 W, t6 g! m; }    fn drop(&amp;mut self) {
3 F, S! t! m8 E* ]        for _ in 0..self.max_workers {, ^- b3 b# O2 m+ t- S
            self.sender.send(Message::ByeBye).unwrap();; S& M' P7 T9 |. ?, z
        }: F9 M5 U' Y2 \9 n/ e
        for w in self.workers {4 s. E! N$ g! z' l& g
            if let Some(t) = w.t.take() {) F2 A; ]$ y; V. u4 `
                t.join().unwrap();: z$ [8 e# d, o% E. @2 j9 g
            }
4 h: V; F: ]" f" d! a        }- g& Q: [; |% ~
    }
, X; M7 i, ]$ i" l}8 v- b$ {1 o+ P* X. }

% ]; A$ ?# C6 c% c8 {# K1 f$ u9 Y! H: q  ]0 c( B" P  t
#[cfg(test)]
3 A4 w$ Z. D' S: Smod tests {
9 ~' P6 Z0 N3 x% J7 q5 h8 q' U    use super::*;
; ^! @& ^& i9 R7 C    #[test]: m: P9 J0 b9 a( |/ W' h' I7 V
    fn it_works() {
, D! I7 k/ F- g        let p = Pool::new(4);* I" I$ e# T$ D7 H
        p.execute(|| println!("do new job1"));
5 P1 z- K5 q5 |0 X6 ^        p.execute(|| println!("do new job2"));
0 g- y8 l% @$ e5 N        p.execute(|| println!("do new job3"));
. t  a. \- _  S, v        p.execute(|| println!("do new job4"));
8 L. w0 F' o: a( G( P/ Y" n    }" K% l* r+ ?( \# n
}+ P6 ]  k& k5 _0 m/ _
</code></pre>
, z5 Y$ r+ l/ Y1 C
& `+ \  F: d+ j) V+ X6 b2 n
回复

使用道具 举报

懒得打字嘛,点击右侧快捷回复 【右侧内容,后台自定义】
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

手机版|飞雪团队

GMT+8, 2026-9-2 23:52 , Processed in 0.258195 second(s), 21 queries , Gzip On.

Powered by Discuz! X3.4

Copyright © 2001-2021, Tencent Cloud.

快速回复 返回顶部 返回列表