우아한 종료와 정리

우아한 종료와 정리 (Graceful Shutdown and Cleanup)

21-20 목록의 코드는 우리가 의도한 대로 스레드 풀을 통해 요청을 비동기적으로 처리하고 있어요. 다만 직접적인 방식으로 사용하지 않는 workers, id, thread 필드에 대한 경고가 몇 개 뜨는데, 이는 우리가 아무것도 정리(cleanup)하지 않고 있다는 사실을 상기시켜 줘요. 덜 우아한 <kbd>ctrl</kbd>-<kbd>C</kbd> 방식으로 메인 스레드를 멈추면, 요청을 처리하는 도중이더라도 다른 모든 스레드도 즉시 중지돼요.

그럼 이제 Drop 트레이트를 구현해서 풀의 각 스레드에 join을 호출하고, 닫히기 전에 처리 중인 요청을 끝마치게 할게요. 그리고 나서 스레드들에게 새 요청 수락을 멈추고 종료하라고 알리는 방법도 구현할 거예요. 이 코드가 실제로 동작하는 걸 보려고, 서버가 두 개의 요청만 처리한 뒤 스레드 풀을 우아하게 종료하도록 수정할게요.

진행하면서 한 가지 눈여겨볼 점은, 이 모든 것이 클로저 실행을 처리하는 코드 부분에는 영향을 주지 않는다는 거예요. 그래서 여기 모든 내용은 스레드 풀을 async 런타임에 사용하더라도 동일할 겁니다.

출처: The Rust Book

ThreadPoolDrop 트레이트 구현하기

먼저 스레드 풀에 Drop을 구현하는 것부터 시작할게요. 풀이 드롭되면, 스레드들이 각자의 작업을 끝마치는지 확인하기 위해 모두 join해야 돼요. 21-22 목록은 Drop 구현의 첫 시도인데, 이 코드는 아직 완벽하게 동작하진 않아요.

use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: mpsc::Sender<Job>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool { workers, sender }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        for worker in &mut self.workers {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

struct Worker {
    id: usize,
    thread: thread::JoinHandle<()>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let job = receiver.lock().unwrap().recv().unwrap();

                println!("Worker {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}

먼저 스레드 풀의 workers 각각을 루프로 순회해요. self가 가변 참조이기 때문에, 그리고 worker를 변경해야 하기 때문에 &mut를 사용해요. 각 worker에 대해 이 특정 Worker 인스턴스가 종료되고 있다는 메시지를 출력하고, 그 스레드에 join을 호출해요. join 호출이 실패하면 unwrap으로 Rust가 패닉을 일으켜 우아하지 않은 종료로 빠지게 해요.

이 코드를 컴파일하면 다음과 같은 에러가 나요:

<console_highlight>

$ cargo check
    Checking hello v0.1.0 (file:///projects/hello)
error[E0507]: cannot move out of `worker.thread` which is behind a mutable reference
  --> src/lib.rs:52:13
   |
52 |             worker.thread.join().unwrap();
   |             ^^^^^^^^^^^^^ ------ `worker.thread` moved due to this method call
   |             |
   |             move occurs because `worker.thread` has type `JoinHandle<()>`, which does not implement the `Copy` trait
   |
note: `JoinHandle::<T>::join` takes ownership of the receiver `self`, which moves `worker.thread`
  --> /rustc/1159e78c4747b02ef996e55082b704c09b970588/library/std/src/thread/mod.rs:1921:17

For more information about this error, try `rustc --explain E0507`.
error: could not compile `hello` (lib) due to 1 previous error

</console_highlight>

에러는 우리가 각 worker의 가변 대여(mutable borrow)만 갖고 있고 join은 인자의 소유권을 가져가기 때문에, join을 호출할 수 없다고 알려줘요. 이 문제를 해결하려면, join이 스레드를 소비할 수 있도록 thread를 소유한 Worker 인스턴스 밖으로 스레드를 옮겨야 해요. 한 가지 방법은 18-15 목록에서 취했던 것과 같은 접근 방식을 쓰는 거예요. WorkerOption<thread::JoinHandle<()>>을 보관한다면, Optiontake 메서드를 호출해 Some 변형에서 값을 빼내고 그 자리에 None 변형을 남길 수 있어요. 다시 말해, 실행 중인 WorkerthreadSome 변형을 갖고 있고, Worker를 정리하고 싶을 때 SomeNone으로 바꿔서 Worker가 실행할 스레드가 없게 만드는 거죠.

하지만 이 경우는 오직 Worker를 드롭할 때만 발생해요. 그 대가로, worker.thread에 접근하는 모든 곳에서 Option<thread::JoinHandle<()>>을 다뤄야 하죠. 관용적인 Rust는 Option을 꽤 많이 사용하지만, 항상 존재할 거라고 아는 무언가를 이렇게 임시방편으로 Option에 감싸는 상황을 발견하면, 코드를 더 깔끔하고 에러에 덜 취약하게 만들 대안을 찾아보는 게 좋아요.

이 경우 더 나은 대안이 존재해요: Vec::drain 메서드예요. 이 메서드는 벡터에서 어떤 항목을 제거할지 지정하는 범위 인자를 받고, 그 항목들의 반복자를 반환해요. .. 범위 문법을 넘기면 벡터에서 모든 값을 제거해요.

그래서 ThreadPooldrop 구현을 다음과 같이 업데이트해야 해요:

#![allow(unused)]
fn main() {
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: mpsc::Sender<Job>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool { workers, sender }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

struct Worker {
    id: usize,
    thread: thread::JoinHandle<()>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let job = receiver.lock().unwrap().recv().unwrap();

                println!("Worker {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}
}

이렇게 하면 컴파일러 에러가 해결되고, 코드의 다른 어떤 변경도 필요하지 않아요. drop은 패닉이 발생한 상태에서도 호출될 수 있기 때문에, unwrap 역시 패닉을 일으켜 이중 패닉(double panic)을 만들 수 있고, 그러면 프로그램이 즉시 중단되고 진행 중이던 정리가 끝나버린다는 점을 명심하세요. 예제 프로그램에서는 괜찮지만, 프로덕션 코드에서는 권장되지 않아요.

스레드들이 작업을 듣는 걸 멈추도록 신호하기

지금까지의 모든 변경으로 코드는 경고 없이 컴파일돼요. 하지만 안타깝게도 이 코드는 아직 우리가 원하는 대로 동작하지 않아요. 핵심은 Worker 인스턴스들의 스레드가 실행하는 클로저의 로직이에요. 지금은 join을 호출하지만, 스레드들이 작업을 찾으러 loop를 영원히 돌기 때문에 그건 스레드를 종료시키지 못해요. 현재의 drop 구현으로 ThreadPool을 드롭하려 하면, 메인 스레드는 첫 스레드가 끝나길 기다리며 영원히 차단돼 버려요.

이 문제를 고치려면 ThreadPooldrop 구현을 바꾸고, 그다음 Worker의 루프를 바꿔야 해요.

먼저, 스레드들이 끝나길 기다리기 전에 sender를 명시적으로 드롭하도록 ThreadPooldrop 구현을 바꿀게요. 21-23 목록은 sender를 명시적으로 드롭하도록 변경한 ThreadPool을 보여줘요. 스레드와는 달리, 여기서는 Option::takesenderThreadPool에서 옮겨내기 위해 Option을 사용할 필요가 있어요.

use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}
// --snip--

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        // --snip--

        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool {
            workers,
            sender: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

struct Worker {
    id: usize,
    thread: thread::JoinHandle<()>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let job = receiver.lock().unwrap().recv().unwrap();

                println!("Worker {id} got a job; executing.");

                job();
            }
        });

        Worker { id, thread }
    }
}

sender를 드롭하면 채널이 닫혀서, 더 이상 메시지가 보내지지 않는다는 걸 알려줘요. 그렇게 되면 Worker 인스턴스들이 무한 루프에서 수행하는 모든 recv 호출은 에러를 반환하게 돼요. 21-24 목록에서는 Worker 루프를 바꿔 그런 경우에 우아하게 루프를 빠져나가게 해요. 즉 ThreadPooldrop 구현이 그들에 대해 join을 호출할 때 스레드들이 끝나게 된다는 뜻이에요.

use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool {
            workers,
            sender: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in self.workers.drain(..) {
            println!("Shutting down worker {}", worker.id);

            worker.thread.join().unwrap();
        }
    }
}

struct Worker {
    id: usize,
    thread: thread::JoinHandle<()>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let message = receiver.lock().unwrap().recv();

                match message {
                    Ok(job) => {
                        println!("Worker {id} got a job; executing.");

                        job();
                    }
                    Err(_) => {
                        println!("Worker {id} disconnected; shutting down.");
                        break;
                    }
                }
            }
        });

        Worker { id, thread }
    }
}

이 코드가 실제로 동작하는 걸 보려면, main을 수정해서 두 개의 요청만 처리한 뒤 서버를 우아하게 종료하도록 만들어 볼게요. 21-25 목록을 봐주세요.

use hello::ThreadPool;
use std::{
    fs,
    io::{BufReader, prelude::*},
    net::{TcpListener, TcpStream},
    thread,
    time::Duration,
};

fn main() {
    let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
    let pool = ThreadPool::new(4);

    for stream in listener.incoming().take(2) {
        let stream = stream.unwrap();

        pool.execute(|| {
            handle_connection(stream);
        });
    }

    println!("Shutting down.");
}

fn handle_connection(mut stream: TcpStream) {
    let buf_reader = BufReader::new(&stream);
    let request_line = buf_reader.lines().next().unwrap().unwrap();

    let (status_line, filename) = match &request_line[..] {
        "GET / HTTP/1.1" => ("HTTP/1.1 200 OK", "hello.html"),
        "GET /sleep HTTP/1.1" => {
            thread::sleep(Duration::from_secs(5));
            ("HTTP/1.1 200 OK", "hello.html")
        }
        _ => ("HTTP/1.1 404 NOT FOUND", "404.html"),
    };

    let contents = fs::read_to_string(filename).unwrap();
    let length = contents.len();

    let response =
        format!("{status_line}\r\nContent-Length: {length}\r\n\r\n{contents}");

    stream.write_all(response.as_bytes()).unwrap();
}

실제 현실 세계의 웹 서버가 두 개의 요청만 처리하고 종료하길 원하진 않을 거예요. 이 코드는 우아한 종료와 정리가 제대로 동작한다는 걸 보여주기 위한 것일 뿐이에요.

take 메서드는 Iterator 트레이트에 정의되어 있고, 반복을 최대 처음 두 항목으로 제한해요. ThreadPoolmain의 끝에서 스코프를 벗어나고, drop 구현이 실행될 거예요.

cargo run으로 서버를 시작하고 요청 세 개를 보내 보세요. 세 번째 요청은 에러가 나야 하고, 터미널에서 다음과 비슷한 출력을 볼 수 있을 거예요:

<console_highlight>

$ cargo run
   Compiling hello v0.1.0 (file:///projects/hello)
    Finished `dev` profile [unoptimized + debuginfo] target(s) in 0.41s
     Running `target/debug/hello`
Worker 0 got a job; executing.
Shutting down.
Shutting down worker 0
Worker 3 got a job; executing.
Worker 1 disconnected; shutting down.
Worker 2 disconnected; shutting down.
Worker 3 disconnected; shutting down.
Worker 0 disconnected; shutting down.
Shutting down worker 1
Shutting down worker 2
Shutting down worker 3

</console_highlight>

Worker ID와 메시지가 출력되는 순서는 다를 수 있어요. 메시지들에서 이 코드가 어떻게 동작하는지 볼 수 있어요. Worker 인스턴스 0과 3이 첫 두 요청을 받았어요. 서버는 두 번째 연결 이후에 새 연결 수락을 멈췄고, ThreadPoolDrop 구현은 Worker 3이 자신의 작업을 시작하기도 전에 실행되기 시작해요. sender를 드롭하면 모든 Worker 인스턴스의 연결이 끊기고 종료하도록 알려줘요. Worker 인스턴스들은 각자 연결이 끊길 때 메시지를 출력하고, 그다음 스레드 풀이 join을 호출해 각 Worker 스레드가 끝나길 기다려요.

이 특정 실행에서 흥미로운 점 하나를 눈여겨보세요. ThreadPoolsender를 드롭했고, 어떤 Worker도 에러를 받기 전에 우리는 Worker 0join을 시도했어요. Worker 0은 아직 recv에서 에러를 받지 않았기 때문에, 메인 스레드는 Worker 0이 끝나길 기다리며 차단됐어요. 그 사이 Worker 3이 작업을 받았고, 그다음 모든 스레드가 에러를 받았어요. Worker 0이 끝나면, 메인 스레드는 나머지 Worker 인스턴스들이 끝나길 기다렸어요. 그 시점에 그들은 모두 루프를 빠져나와 멈춰 있었죠.

축하해요! 이제 프로젝트를 완성했어요. 스레드 풀을 사용해 비동기적으로 응답하는 기본 웹 서버를 만들었답니다. 서버를 스레드 풀의 모든 스레드를 정리하는 우아한 종료를 수행할 수 있어요.

참고를 위해 전체 코드는 다음과 같아요:

use hello::ThreadPool;
use std::{
    fs,
    io::{BufReader, prelude::*},
    net::{TcpListener, TcpStream},
    thread,
    time::Duration,
};

fn main() {
    let listener = TcpListener::bind("127.0.0.1:7878").unwrap();
    let pool = ThreadPool::new(4);

    for stream in listener.incoming().take(2) {
        let stream = stream.unwrap();

        pool.execute(|| {
            handle_connection(stream);
        });
    }

    println!("Shutting down.");
}

fn handle_connection(mut stream: TcpStream) {
    let buf_reader = BufReader::new(&stream);
    let request_line = buf_reader.lines().next().unwrap().unwrap();

    let (status_line, filename) = match &request_line[..] {
        "GET / HTTP/1.1" => ("HTTP/1.1 200 OK", "hello.html"),
        "GET /sleep HTTP/1.1" => {
            thread::sleep(Duration::from_secs(5));
            ("HTTP/1.1 200 OK", "hello.html")
        }
        _ => ("HTTP/1.1 404 NOT FOUND", "404.html"),
    };

    let contents = fs::read_to_string(filename).unwrap();
    let length = contents.len();

    let response =
        format!("{status_line}\r\nContent-Length: {length}\r\n\r\n{contents}");

    stream.write_all(response.as_bytes()).unwrap();
}
use std::{
    sync::{Arc, Mutex, mpsc},
    thread,
};

pub struct ThreadPool {
    workers: Vec<Worker>,
    sender: Option<mpsc::Sender<Job>>,
}

type Job = Box<dyn FnOnce() + Send + 'static>;

impl ThreadPool {
    /// Create a new ThreadPool.
    ///
    /// The size is the number of threads in the pool.
    ///
    /// # Panics
    ///
    /// The `new` function will panic if the size is zero.
    pub fn new(size: usize) -> ThreadPool {
        assert!(size > 0);

        let (sender, receiver) = mpsc::channel();

        let receiver = Arc::new(Mutex::new(receiver));

        let mut workers = Vec::with_capacity(size);

        for id in 0..size {
            workers.push(Worker::new(id, Arc::clone(&receiver)));
        }

        ThreadPool {
            workers,
            sender: Some(sender),
        }
    }

    pub fn execute<F>(&self, f: F)
    where
        F: FnOnce() + Send + 'static,
    {
        let job = Box::new(f);

        self.sender.as_ref().unwrap().send(job).unwrap();
    }
}

impl Drop for ThreadPool {
    fn drop(&mut self) {
        drop(self.sender.take());

        for worker in &mut self.workers {
            println!("Shutting down worker {}", worker.id);

            if let Some(thread) = worker.thread.take() {
                thread.join().unwrap();
            }
        }
    }
}

struct Worker {
    id: usize,
    thread: Option<thread::JoinHandle<()>>,
}

impl Worker {
    fn new(id: usize, receiver: Arc<Mutex<mpsc::Receiver<Job>>>) -> Worker {
        let thread = thread::spawn(move || {
            loop {
                let message = receiver.lock().unwrap().recv();

                match message {
                    Ok(job) => {
                        println!("Worker {id} got a job; executing.");

                        job();
                    }
                    Err(_) => {
                        println!("Worker {id} disconnected; shutting down.");
                        break;
                    }
                }
            }
        });

        Worker {
            id,
            thread: Some(thread),
        }
    }
}

여기서 더 할 수 있는 것들이 많아요! 이 프로젝트를 계속 개선하고 싶다면 몇 가지 아이디어를 알려드릴게요:

  • ThreadPool과 그 공개 메서드들에 문서를 더 추가하세요.
  • 라이브러리의 기능에 대한 테스트를 추가하세요.
  • unwrap 호출을 더 견고한 에러 처리로 바꾸세요.
  • 웹 요청을 처리하는 것 외의 작업을 수행하는 데 ThreadPool을 사용해 보세요.
  • crates.io에서 스레드 풀 크레이트를 찾아 그 크레이트를 대신 사용하는 비슷한 웹 서버를 구현해 보세요. 그다음 그 API와 견고성을 우리가 구현한 스레드 풀과 비교해 보세요.

요약 (Summary)

수고 많았어요! 이 책의 끝까지 왔네요! 이 Rust 여행에 함께해 줘서 고마워요. 이제 여러분의 Rust 프로젝트를 직접 구현하고 다른 사람들의 프로젝트를 도울 준비가 됐어요. Rust 여정에서 맞닥뜨리는 어떤 어려움에든 기꺼이 도와줄 따뜻한 Rustacean 커뮤니티가 기다리고 있다는 걸 잊지 마세요.

더 알아보기 (Learn more)