Concurrency
Comprehension
Concurrency/parallelism can use message passing or shared state
Thread
Concurrency requires attention to:
data races
Race conditions, where threads are accessing data or resources in an inconsistent order
deadlocks
Deadlocks, where two threads are waiting for each other to finish using a resource the other thread has, preventing both threads from continuing
hard-to-reproduce bugs
Bugs that happen only in certain situations and are hard to reproduce and fix reliably
A green thread is a language-level thread, not exactly the same as an OS thread
use std::thread;
use std::time::Duration;
fn main() {
let handle = thread::spawn(|| { // explicitly declare move if this closure needs arguments
for i in 1..10 {
println!("hi number {} from the spawned thread!", i);
thread::sleep(Duration::from_millis(1)); // the OS may switch to another thread here
}
});
handle.join().unwrap(); // use join to wait for all threads (otherwise later code does not run)
for i in 1..5 {
println!("hi number {} from the main thread!", i);
thread::sleep(Duration::from_millis(1));
}
}
Message-passing
Basic channel syntax
use std::sync::mpsc; // multiple producer single consumer fn main() { let (tx, rx) = mpsc::channel(); // transmitter, receiver thread::spawn(move || { let val = String::from("hi"); tx.send(val).unwrap(); // val ownership is moved here }); let received = rx.recv().unwrap(); // try_recv does not block println!("Got: {}", received); }Multi producer/transmitter && Receive multi message
use std::sync::mpsc; use std::thread; use std::time::Duration; fn main() { let (tx, rx) = mpsc::channel(); // the later closure moves tx ownership // clone tx into tx1 first for multiple producers let tx1 = tx.clone(); thread::spawn(move || { let vals = vec![ String::from("hi"), String::from("from"), String::from("the"), String::from("thread"), ]; for val in vals { tx1.send(val).unwrap(); thread::sleep(Duration::from_secs(1)); } }); thread::spawn(move || { let vals = vec![ String::from("more"), String::from("messages"), String::from("for"), String::from("you"), ]; for val in vals { tx.send(val).unwrap(); thread::sleep(Duration::from_secs(1)); } }); // note that this uses a for-in loop for received in rx { println!("Got: {}", received); } }
Shared-state
A mutex (mutual exclusion) allows only one thread to access particular data at a time
A thread must acquire a lock to access data in a mutex. The lock tracks who has exclusive access, so a mutex protects its data through a locking system.
Mutexes are difficult to use because:
you must try to acquire the lock before using the data
You must attempt to acquire the lock before using the data.
you must release the lock when finished so other threads can acquire it
When you’re done with the data that the mutex guards, you must unlock the data so other threads can acquire the lock.
Managing mutexes is difficult, so many prefer channels. Rust’s type system and ownership rules prevent incorrect lock acquisition or release
Basic mutex syntax
use std::sync::Mutex; fn main() { let m = Mutex::new(5); { let mut num = m.lock().unwrap(); // call lock to acquire the lock *num = 6; } println!("m = {:?}", m); }Mutex<T>is a smart pointer. More precisely,lockreturns aLockResult, whoseunwrapreturns the smart pointerMutexGuard.MutexGuardimplementsDerefto point to the inner data andDropto release the lock automatically when it leaves scope. Thus we do not risk forgetting to release the lock or blocking other threads, because release is automatic.Sharing a mutex across threads
Atomic Reference Counted Type,Think of it as an atomic
Rc<T>use std::sync::{Arc, Mutex}; use std::thread; fn main() { let counter = Arc::new(Mutex::new(0)); let mut handles = vec![]; for _ in 0..10 { let counter = Arc::clone(&counter); let handle = thread::spawn(move || { let mut num = counter.lock().unwrap(); *num += 1; }); handles.push(handle); } for handle in handles { handle.join().unwrap(); } println!("Result: {}", *counter.lock().unwrap()); }
Send & Sync trait
Types implementing Send can transfer ownership between threads
Types implementing Sync can be accessed by multiple threads
Thus, Sync includes the behavior of Send
Origin
…
Atomic Reference Counting with Arc<T>
Fortunately, Arc<T> is a type like Rc<T> that is safe to use in concurrent situations. The a stands for atomic, meaning it’s an atomically reference counted type. Atomics are an additional kind of concurrency primitive that we won’t cover in detail here: see the standard library documentation for [std::sync::atomic](https://doc.rust-lang.org/std/sync/atomic/index.html) for more details. At this point, you just need to know that atomics work like primitive types but are safe to share across threads.
You might then wonder why all primitive types aren’t atomic and why standard library types aren’t implemented to use Arc<T> by default. The reason is that thread safety comes with a performance penalty that you only want to pay when you really need to. If you’re just performing operations on values within a single thread, your code can run faster if it doesn’t have to enforce the guarantees atomics provide.
Let’s return to our example: Arc<T> and Rc<T> have the same API, so we fix our program by changing the use line, the call to new, and the call to clone. The code in Listing 16-15 will finally compile and run:
…