Rust - Basic - 15 - Concurrency

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, lock returns a LockResult, whose unwrap returns the smart pointer MutexGuard. MutexGuard implements Deref to point to the inner data and Drop to 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

https://doc.rust-lang.org/book/ch16-00-concurrency.html

…

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:

…