Rust Concurrent Programming

Safe and efficient handling of concurrency is one of the purposes for which Rust was created, mainly addressing the server's ability to withstand high load.

The concept of concurrency refers to different parts of a program executing independently, which is easily confused with the concept of parallelism; parallelism emphasizes "simultaneous execution."

Concurrency often results in parallelism.

This chapter discusses programming concepts and details related to concurrency.

Threads

A thread is a part of a program that runs independently.

The difference between a thread and a process is that a thread is a concept within a program; a program often executes within a process.

In an environment with an operating system, processes are often alternately scheduled for execution, while threads are scheduled by the program within the process.

Since thread concurrency is very likely to result in parallelism, deadlock and delay errors encountered in parallel execution often appear in programs with concurrency mechanisms.

To solve these problems, many other languages (such as Java and C#) use special runtime software to coordinate resources, but this undoubtedly greatly reduces program execution efficiency.

C/C++ support multithreading at the lowest level of the operating system, but the language itself and its compilers lack the ability to detect and avoid parallel errors, which puts great pressure on developers; developers need to spend a lot of effort avoiding errors.

Rust does not rely on a runtime environment, similar to C/C++.

But Rust has designed, within the language itself, mechanisms including the ownership system to eliminate the most common errors at the compilation stage as much as possible, which other languages do not have.

But this does not mean we can be careless when programming; to date, problems caused by concurrency have not been completely resolved in the public domain, and errors can still occur. Be as careful as possible when programming concurrently!

In Rust, new threads are created through the std::thread::spawn function:

Example

use std::thread;
use std::time::Duration;

fn spawn_function() {
    for i in 0..5 {
        println!("spawned thread print {}", i);
        thread::sleep(Duration::from_millis(1));
    }
}

fn main() {
    thread::spawn(spawn_function);

    for i in 0..3 {
        println!("main thread print {}", i);
        thread::sleep(Duration::from_millis(1));
    }
}

Output:

main thread print 0
spawned thread print 0
main thread print 1
spawned thread print 1
main thread print 2
spawned thread print 2

The order of this result may change in some cases, but generally this is how it is printed.

This program has a child thread that prints 5 lines of text, and the main thread prints 3 lines of text. But clearly, when the main thread ends, the spawn thread also ends without completing all the printing.

The parameter of the std::thread::spawn function is a function that takes no parameters, but the above approach is not recommended. We can use closures to pass a function as an argument:

Example

use std::thread;
use std::time::Duration;

fn main() {
    thread::spawn(|| {
        for i in 0..5 {
            println!("spawned thread print {}", i);
            thread::sleep(Duration::from_millis(1));
        }
    });

    for i in 0..3 {
        println!("main thread print {}", i);
        thread::sleep(Duration::from_millis(1));
    }
}

A closure is an anonymous function that can be stored in a variable or passed as an argument to other functions. Closures are equivalent to Lambda expressions in Rust, with the following format:

|参数1, 参数2, ...| -> 返回值类型 {
    // 函数体
}

For example:

Example

fn main() {
    let inc = |num: i32| -> i32 {
        num + 1
    };
    println!("inc(5) = {}", inc(5));
}

Output:

inc(5) = 6

Closures can omit type declarations and use Rust's automatic type inference mechanism:

Example

fn main() {
    let inc = |num| {
        num + 1
    };
    println!("inc(5) = {}", inc(5));
}

The result is unchanged.

join method

Example

use std::thread;
use std::time::Duration;

fn main() {
    let handle = thread::spawn(|| {
        for i in 0..5 {
            println!("spawned thread print {}", i);
            thread::sleep(Duration::from_millis(1));
        }
    });

    for i in 0..3 {
        println!("main thread print {}", i);
        thread::sleep(Duration::from_millis(1));
    }

    handle.join().unwrap();
}

Output:

main thread print 0 
spawned thread print 0 
spawned thread print 1 
main thread print 1 
spawned thread print 2 
main thread print 2 
spawned thread print 3 
spawned thread print 4

The join method causes the program to stop running only after the child thread has finished.

move forces ownership transfer

This is a commonly encountered situation:

Example

use std::thread;

fn main() {
    let s = "hello";
   
    let handle = thread::spawn(|| {
        println!("{}", s);
    });

    handle.join().unwrap();
}

Trying to use the current function's resources in a child thread is definitely wrong! Because the ownership mechanism prohibits this dangerous situation; it would break the certainty of resource destruction that the ownership mechanism provides. We can use the move keyword of closures to handle this:

Example

use std::thread;

fn main() {
    let s = "hello";
   
    let handle = thread::spawn(move || {
        println!("{}", s);
    });

    handle.join().unwrap();
}

Message Passing

A primary tool in Rust for implementing message-passing concurrency is the channel, which consists of two parts: a transmitter and a receiver.

std::sync::mpsc contains the methods for message passing:

Example

use std::thread;
use std::sync::mpsc;

fn main() {
    let (tx, rx) = mpsc::channel();

    thread::spawn(move || {
        let val = String::from("hi");
        tx.send(val).unwrap();
    });

    let received = rx.recv().unwrap();
    println!("Got: {}", received);
}

Output:

Got: hi

The child thread obtained the sender tx from the main thread, called its send method to send a string, and then the main thread received it through the corresponding receiver rx.

Other extensions