Rust Бесстрашная конкурентность: потоки, каналы, Send и Sync, разделяемое состояние
0%

Бесстрашная конкурентность: потоки, каналы, Send и Sync, разделяемое состояние

Бесстрашная конкурентность: потоки, каналы, Send и Sync, разделяемое состояние

«Fearless concurrency» — единственный маркетинговый слоган Rust, который является техническим утверждением с доказательством. Он не говорит «многопоточность станет простой». Он говорит ровно одно: в безопасном подмножестве языка гонка данных невыразима — программа с гонкой данных не компилируется.

Это не новый механизм. Всё уже собрано в владении и заимствовании: правило «либо один изменяющий доступ, либо сколько угодно читающих» — это и есть определение отсутствия гонки, только сформулированное без слова «поток». Конкурентность в Rust достраивает над ним два трейта-маркера, Send и Sync, и на этом теория заканчивается. Дальше — инженерия: какие примитивы выбрать, сколько они стоят, обо что вы всё равно споткнётесь.

Задача: один счётчик и восемь потоков

Начнём с классики. Восемь потоков увеличивают общий счётчик, каждый по пять миллионов раз. На C это пишется в три строки и работает неправильно:

// C: компилируется, запускается, даёт неверный ответ — молча
static long counter = 0;                  // общая переменная
void *worker(void *_) {
    for (int i = 0; i < 5000000; i++) counter++;   // read-modify-write без синхронизации
    return NULL;
}
// 8 потоков → ожидаем 40000000, получаем, например, 9137422. Каждый запуск — своё число.

Почему так, подробно разобрано в «Потоки и синхронизация»: counter++ — это три инструкции (загрузка, инкремент, запись), и между ними влезает другой поток. Хуже другое: это ещё и неопределённое поведение, а не просто «потерянные инкременты». Компилятор имеет право закэшировать counter в регистре и вынести его из цикла — тогда результат станет вообще произвольным («Неопределённое поведение»).

Дадим точное определение, потому что дальше всё держится на нём. Гонка данных — это одновременное выполнение трёх условий:

  1. два или более потока обращаются к одной ячейке памяти;
  2. хотя бы одно обращение — запись;
  3. обращения не упорядочены отношением synchronizes-with (нет мьютекса, атомика, join — ничего).

Уберите любое одно условие — гонки нет. Rust убирает первое и второе одновременно через систему типов: если к данным есть изменяющий доступ, то он единственный; если доступов много, все они только читают. Гонка требует ровно того, что запрещено правилом заимствования.

Что «бесстрашная конкурентность» обещает и что не обещает

Это самый важный раздел главы, и его обычно пропускают. Компилятор гарантирует:

  • нет гонок данных — ни в одной строке безопасного кода, включая код через Arc, каналы, rayon, чужие крейты;
  • нет use-after-free между потоками — поток не может пережить данные, которые заимствовал ('static или scope);
  • нет «забыл взять замок» — данные физически лежат внутри Mutex, другого пути к ним нет;
  • нет «забыл отпустить замок»MutexGuard освобождается в drop, в том числе при панике;
  • нет неатомарного счётчика ссылок между потокамиRc просто не пройдёт в поток.

Компилятор не гарантирует:

  • отсутствие взаимоблокировок. Два Mutex, захваченные в разном порядке двумя потоками, — компилируется, проходит clippy, висит в проде.
  • отсутствие логических состояний гонки (race condition в широком смысле). if !map.contains_key(k) { map.insert(k, v) } под двумя последовательными lock() — корректно с точки зрения памяти и неверно по смыслу.
  • правильность порядков памяти в атомиках. Relaxed там, где нужен Acquire, не даёт UB, но даёт неверные результаты.
  • отсутствие голодания, живой блокировки, инверсии приоритетов.
  • отсутствие утечек и OOM. Неограниченный канал, у которого производитель быстрее потребителя, съест всю память («Обработка ошибок» про то, почему это надо считать штатным сценарием).
  • производительность. Восемь потоков на одном мьютексе работают медленнее одного потока — и это самая частая реальная проблема, а не гонки.

Формулировка, которую стоит запомнить: Rust убирает класс ошибок, который нельзя отладить (нелокальный, невоспроизводимый, зависящий от планировщика), и оставляет класс ошибок, который отлаживается (висим — снимем стеки всех потоков и увидим, кто кого ждёт).

Поток как ресурс операционной системы

std::thread — это тонкая обёртка над потоками ОС: pthread_create на Unix, CreateThread на Windows. Никаких зелёных потоков в стандартной библиотеке нет (были до 1.0 и были убраны). Один thread::spawn — один поток ядра со своим стеком, которого планировщик ОС видит наравне с остальными («Процессы и планирование»).

use std::thread;

fn main() {
    let mut handles = Vec::new();

    for id in 0..4 {
        // spawn возвращает JoinHandle<T>, где T — тип возвращаемого значения замыкания
        handles.push(thread::spawn(move || {
            let sum: u64 = (1..=1000u64).map(|x| x * id).sum();
            format!("поток {id} насчитал {sum}")
        }));
    }

    for h in handles {
        // join() блокирует до завершения и ВОЗВРАЩАЕТ результат:
        // Ok(значение) либо Err(payload паники)
        match h.join() {
            Ok(msg) => println!("{msg}"),
            Err(_)  => eprintln!("поток паниковал"),
        }
    }
}

Три вещи, которые надо знать про эту обёртку сразу.

Поток дорог. Создание — десятки микросекунд, стек по умолчанию 8 МиБ виртуальной памяти (резервируется, коммитится по факту использования). Тысяча потоков — это тысяча стеков и тысяча кандидатов на переключение контекста. Отсюда две стратегии: пул потоков фиксированного размера под CPU-работу и асинхронность под ожидание ввода-вывода («Асинхронный Rust»). Разумный размер пула — std::thread::available_parallelism(), а не константа 16:

let n = std::thread::available_parallelism()
    .map(|v| v.get())
    .unwrap_or(1);   // учитывает cgroup-лимиты контейнера, в отличие от num_cpus

Потоки без join умирают вместе с main. Это неочевидно и ловит всех:

use std::{thread, time::Duration};

fn main() {
    thread::spawn(|| {
        thread::sleep(Duration::from_millis(50));
        println!("я никогда не напечатаюсь");   // main уже завершил процесс
    });
    println!("main закончился");
}
// вывод: только «main закончился»

Никакого предупреждения не будет: JoinHandle не помечен #[must_use]. В проде это выглядит как «половина логов иногда теряется при завершении». Правило: либо join, либо явно осознанный «пожизненный» демон-поток.

Даём потокам имена. Это единственный способ понять, что происходит, глядя в gdb, perf или логи:

let h = std::thread::Builder::new()
    .name("ingest-worker".to_string())     // видно в gdb `info threads` и в панике
    .stack_size(2 * 1024 * 1024)          // 2 МиБ вместо 8, если потоков много
    .spawn(|| { /* работа */ })
    .expect("не удалось создать поток");  // spawn через Builder возвращает io::Result

Первая ошибка компилятора: 'static

Попробуем передать в поток ссылку на локальные данные:

use std::thread;

fn main() {
    let data = vec![1, 2, 3];
    let h = thread::spawn(|| println!("{data:?}"));   // не собирается
    h.join().unwrap();
}
error[E0373]: closure may outlive the current function, but it borrows `data`,
              which is owned by the current function
 --> src/main.rs:5:27
  |
5 |     let h = thread::spawn(|| println!("{data:?}"));
  |                           ^^                ---- `data` is borrowed here
  |                           |
  |                           may outlive borrowed value `data`
  |
note: function requires argument type to outlive `'static`
help: to force the closure to take ownership of `data` (and any other referenced
      variables), use the `move` keyword
  |
5 |     let h = thread::spawn(move || println!("{data:?}"));
  |                           ++++

Сообщение стоит разобрать целиком, потому что оно объясняет модель. Сигнатура thread::spawn требует F: FnOnce() -> T + Send + 'static. Ограничение 'static означает не «живёт всю программу», а «не содержит заимствований с ограниченным временем жизни» («Времена жизни»). Причина: компилятор не знает, что вы вызовете join() — вы могли уронить JoinHandle, и тогда поток пережил бы main. Ссылка на стековую переменную стала бы висячей.

move решает задачу переносом владения. Но что делать, если данные нужны и после потоков, а клонировать их дорого?

Область вместо 'static: scoped threads

С Rust 1.63 в стандартной библиотеке есть thread::scope. Он даёт компилятору то, чего не хватало: гарантию, что все порождённые потоки завершатся до выхода из области, — а значит, заимствовать стек можно.

use std::thread;

fn main() {
    let mut data = vec![5u64, 1, 9, 3, 7, 2, 8, 4];
    let mid = data.len() / 2;
    let (left, right) = data.split_at_mut(mid);   // два непересекающихся &mut

    thread::scope(|s| {
        s.spawn(|| left.sort_unstable());         // заимствуем стек — без move и без Arc
        s.spawn(|| right.sort_unstable());
    });   // здесь scope дожидается обоих потоков: неявный join

    // Оба заимствования закончились — data снова доступна целиком
    println!("{data:?}");   // [1, 3, 5, 9, 2, 4, 7, 8] — две отсортированные половины
}

Обратите внимание на split_at_mut: это единственная строчка, где выражена суть параллелизма по данным. Она превращает один &mut [T] в два непересекающихся, и дальше система типов уже сама не даст двум потокам залезть в одну ячейку. Никакой синхронизации не нужно, потому что нечего синхронизировать.

Правило выбора простое: thread::scope — когда время жизни потоков совпадает с фрагментом кода (параллельная обработка куска данных, fork-join). thread::spawn + Arc — когда поток живёт своей жизнью (фоновый флешер, приёмщик соединений).

Send и Sync: два предложения, из которых следует остальное

Вся многопоточная часть системы типов — это два маркерных трейта без методов:

  • T: Send — значение типа T можно переместить в другой поток. Владение мигрирует, старый поток к нему больше не обращается.
  • T: Sync — на значение типа T можно раздать &T нескольким потокам одновременно. Формально: T: Sync тогда и только тогда, когда &T: Send.

Оба — авто-трейты: компилятор выводит их структурно. Тип Send, если все его поля Send; Sync, если все поля Sync. Вручную писать impl не нужно почти никогда — а когда нужно, это unsafe impl, потому что вы берёте на себя обещание, которое компилятор проверить не может.

Практическая ценность в том, что из этих двух определений выводятся все «почему нельзя»:

Тип Send Sync Причина
u64, String, Vec<T: Send> да да обычные данные, никакого разделения внутри
Rc<T> нет нет счётчик ссылок неатомарный: два потока сделают clone — счётчик испортится, будет double free
Arc<T: Send + Sync> да да счётчик атомарный; но T тоже должен пускать к себе несколько потоков
Cell<T>, RefCell<T> да (если T: Send) нет изменение через &self без синхронизации; переместить целиком можно, разделять — нет
Mutex<T: Send> да да синхронизация внутри, поэтому &Mutex<T> безопасно раздать всем
MutexGuard<'_, T> нет да (если T: Sync) POSIX требует, чтобы мьютекс отпускал тот же поток, что запер
*const T, *mut T нет нет компилятор ничего не знает об инвариантах сырого указателя
&T если T: Sync если T: Sync «дать ссылку в поток» = определение Sync
&mut T если T: Send если T: Sync исключительный доступ = переезд доступа

Как читать E0277 про Send

Это сообщение выглядит пугающе длинным, и именно поэтому его надо разобрать один раз внимательно. Попробуем протащить Rc в поток:

use std::rc::Rc;
use std::thread;

fn main() {
    let shared = Rc::new(vec![1, 2, 3]);
    let cloned = Rc::clone(&shared);
    thread::spawn(move || println!("{}", cloned.len())).join().unwrap();
}
error[E0277]: `Rc<Vec<i32>>` cannot be sent between threads safely
   --> src/main.rs:7:19
    |
7   |     thread::spawn(move || println!("{}", cloned.len())).join().unwrap();
    |     ------------- ^------
    |     |             |
    |     |             within this `{closure@src/main.rs:7:19: 7:26}`
    |     required by a bound introduced by this call
    |
    = help: within `{closure@src/main.rs:7:19: 7:26}`, the trait `Send` is not
            implemented for `Rc<Vec<i32>>`
note: required because it's used within this closure
   --> src/main.rs:7:19
note: required by a bound in `std::thread::spawn`
   --> /rustc/.../library/std/src/thread/mod.rs:731:8
    |
728 | pub fn spawn<F, T>(f: F) -> JoinHandle<T>
    |        ----- required by a bound in this function
...
731 |     F: Send + 'static,
    |        ^^^^ required by this bound in `spawn`

Читается снизу вверх: (1) spawn требует F: Send; (2) F — это ваше замыкание; (3) замыкание не Send, потому что захватило поле типа Rc<Vec<i32>>; (4) Rc не Send. Метка within this {closure@...} — самая полезная: она говорит, что проблема не в самом замыкании, а в том, что внутри. Когда захватов много, ищите строку the trait Send is not implemented for ... — там имя виновника.

Типичное продолжение диалога с компилятором — замена Rc на Arc, после которой возникает второе сообщение, если внутри был RefCell:

error[E0277]: `RefCell<Vec<i32>>` cannot be shared between threads safely
    = help: the trait `Sync` is not implemented for `RefCell<Vec<i32>>`
note: required for `Arc<RefCell<Vec<i32>>>` to implement `Send`
help: consider using `std::sync::Mutex` or `std::sync::RwLock` instead

Обратите внимание: Arc требует от содержимого и Send, и Sync — иначе через клоны Arc два потока получили бы &RefCell, а RefCell считает свой счётчик заимствований неатомарно. Так что «Rc<RefCell<T>> внутри потока» превращается в «Arc<Mutex<T>> между потоками» — это не два разных приёма, а один и тот же с атомарными версиями обеих частей.

Когда unsafe impl всё-таки нужен

Обёртка над сырым указателем (FFI-хендл, самодельная структура) теряет Send/Sync — вернуть их можно только обещанием, а обратный приём, PhantomData<*const ()> в поле, наоборот снимает их, ничего не меняя в раскладке:

struct DeviceHandle(*mut std::ffi::c_void);

// Обещание: библиотека документирует, что хендл работает из любого потока,
// но только из одного одновременно. Значит Send есть, Sync нет.
unsafe impl Send for DeviceHandle {}

Подробности — в «unsafe и FFI». Здесь важно лишь то, что unsafe impl Send — это не «отключить проверку», а «предоставить доказательство самому».

Каналы: передача владения вместо разделения

Самый безопасный способ обойтись без разделяемого состояния — не разделять его. Канал перемещает значение из потока в поток: отправитель теряет владение, получатель приобретает. Синхронизация становится не вашей проблемой, потому что в каждый момент у данных один хозяин. Это модель CSP, как в Go («Конкурентность в Go»), и родственник акторов в Elixir («Конкурентность в Elixir») — только владение здесь проверяется статически, а не обеспечивается копированием сообщений.

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

#[derive(Debug)]
struct Job { id: u32, payload: String }

fn main() {
    // mpsc = multi-producer, single-consumer. Канал НЕОГРАНИЧЕННЫЙ.
    let (tx, rx) = mpsc::channel::<Job>();

    for id in 0..3 {
        let tx = tx.clone();          // каждому потоку — свой Sender
        thread::spawn(move || {
            let job = Job { id, payload: format!("данные-{id}") };
            // send отдаёт владение job; при ошибке ВОЗВРАЩАЕТ его обратно в SendError
            if let Err(e) = tx.send(job) {
                eprintln!("получатель умер, задача не потеряна: {:?}", e.0);
            }
        });
    }
    drop(tx);   // КРИТИЧНО: исходный Sender тоже надо уронить

    // for по Receiver = while let Ok(job) = rx.recv()
    // Цикл закончится, когда умрёт ПОСЛЕДНИЙ Sender.
    for job in rx {
        println!("получено {job:?}");
    }
    println!("все отправители закрыты, обработка завершена");
}

Строчка drop(tx) — грабля номер один во всей главе. Если её убрать, программа напечатает три задачи и зависнет навсегда: в main остался живой Sender, канал не считается закрытым, recv() ждёт четвёртое сообщение. Симптом в проде — «сервис не завершается по SIGTERM».

Обратная сторона того же протокола: send возвращает Err(SendError(value)), если все Receiver уничтожены, — и отдаёт значение назад, чтобы вы могли его залогировать или переложить в другую очередь. Ни один байт не теряется молча.

Ограниченный канал и обратное давление

mpsc::channel не имеет предела. Если производитель быстрее потребителя, очередь растёт до OOM — и это не гипотетика, а один из самых частых способов уронить сервис на Rust. Ограниченный вариант блокирует отправителя:

use std::sync::mpsc::sync_channel;
use std::{thread, time::Duration};

fn main() {
    let (tx, rx) = sync_channel::<u32>(4);   // ёмкость 4; 0 = рандеву без буфера

    let producer = thread::spawn(move || {
        for i in 0..10 {
            tx.send(i).unwrap();             // на 5-м вызове блокируется, пока не заберут
            println!("отправил {i}");
        }
    });

    thread::sleep(Duration::from_millis(100));
    for v in rx { println!("        получил {v}"); }
    producer.join().unwrap();
}

Вывод покажет, что первые четыре отправил появляются мгновенно, а дальше отправитель идёт в ногу с получателем. Это и есть обратное давление: очередь перестаёт быть неявным буфером неограниченного размера. Практическое правило прода — ограниченный канал по умолчанию; неограниченный только там, где вы доказали, что производитель принципиально медленнее.

Когда стандартного mpsc не хватает

  • Нужно несколько получателей (пул воркеров, читающих из одной очереди). Receiver в std помечен !Sync, то есть его нельзя использовать из двух потоков через ссылку. Обходы: Arc<Mutex<Receiver<T>>> (работает, но мьютекс на каждый recv) либо crossbeam-channel, где Receiver клонируется и канал честно mpmc. std::sync::mpmc существует в nightly, но не стабилизирован.
  • Нужно ждать из нескольких каналов сразуcrossbeam_channel::select!, аналог select в Go. В std такого нет.
  • Нужны таймаутыrx.recv_timeout(Duration::from_secs(5)) есть и в std; try_recv() — неблокирующая проверка (не используйте её в цикле без сна, это сжигание процессора).
  • Нужны тик и отменаcrossbeam_channel::tick, after, плюс отдельный «канал завершения»: обычный приём — select! по рабочему каналу и каналу shutdown.

Разделяемое состояние: Arc<Mutex<T>>

Каналы решают не всё. Кэш, метрики, пул соединений, состояние игрового мира — это данные, к которым все обращаются по месту. Здесь нужны две ортогональные вещи: разделяемое владение (кто освободит?) и синхронизация доступа (кто пишет сейчас?). В Rust это два разных типа, и понимание, зачем каждый, — половина дела.

Раскладка Arc, Mutex и MutexGuard в памяти

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

fn main() {
    // Arc — разделяемое владение; Mutex — синхронизация. Ровно один Mutex, восемь ссылок.
    let counter = Arc::new(Mutex::new(0u64));
    let mut handles = Vec::new();

    for _ in 0..8 {
        let counter = Arc::clone(&counter);   // атомарный инкремент счётчика ссылок
        handles.push(thread::spawn(move || {
            for _ in 0..5_000_000 {
                // lock() → LockResult<MutexGuard>. unwrap паникует только при отравлении.
                let mut guard = counter.lock().unwrap();
                *guard += 1;
                // guard уничтожается здесь: замок отпущен, futex_wake разбудит ждущего
            }
        }));
    }
    for h in handles { h.join().unwrap(); }

    println!("итог = {}", *counter.lock().unwrap());   // 40000000 — ВСЕГДА
}

Ключевое отличие от C, pthread_mutex_t и synchronized в Java: Mutex<T> владеет данными. В C мьютекс и защищаемая структура — два независимых объекта, а их связь живёт в комментарии; «забыл взять замок» — обычная находка на ревью, если её вообще заметят. В Rust до T нет пути, кроме lock(), а lock() физически не может вернуть данные, не заперев замок. Правило дисциплины превратилось в свойство типа — это и есть весь фокус.

MutexGuard — RAII-объект: пока он жив, замок заперт; его Drop отпускает замок, в том числе при раскрутке стека после паники. Забыть unlock невозможно. Зато возможна противоположная ошибка — держать guard дольше нужного, о ней ниже.

RwLock, parking_lot, Condvar

RwLock<T> разрешает много читателей или одного писателя: read() и write() возвращают разные guard-типы. Полезен, когда чтений на порядки больше записей и критическая секция не микроскопическая. Иначе Mutex быстрее: у RwLock дороже и захват, и освобождение, а на короткой секции это перекрывает выигрыш от параллельного чтения. Дополнительный риск — голодание писателя при потоке читателей (политика зависит от платформы и не гарантирована документацией).

parking_lot — популярная замена: Mutex в 1 байт вместо 8+, быстрее в неконкурентном случае, честная стратегия ожидания (eventual fairness), нет отравления, есть try_lock_for, а с фичей deadlock_detection — обнаружение взаимоблокировок в рантайме. С Rust 1.62 std на Linux тоже перешёл на futex-реализацию и разрыв сократился, но parking_lot всё ещё выигрывает на многих профилях. Проверяйте на своей нагрузке («Бенчмаркинг»).

Condvar — ответ на «подожди, пока условие станет верным», без сжигания процессора в цикле:

use std::sync::{Condvar, Mutex};

struct Queue {
    items: Mutex<Vec<String>>,
    ready: Condvar,
}

impl Queue {
    fn push(&self, item: String) {
        self.items.lock().unwrap().push(item);
        self.ready.notify_one();          // разбудить одного ждущего
    }

    fn pop(&self) -> String {
        let mut items = self.items.lock().unwrap();
        // wait_while сам делает while-цикл: ложные пробуждения (spurious wakeups) реальны,
        // одного if здесь категорически недостаточно
        items = self.ready.wait_while(items, |v| v.is_empty()).unwrap();
        items.remove(0)
    }
}

wait атомарно отпускает замок и паркует поток, а при пробуждении заново его берёт — поэтому передаётся именно guard. Практический совет: 95% случаев, где хочется Condvar, решаются каналом. Тянитесь к Condvar, когда пишете свой примитив или нужна семантика «разбудить всех» (notify_all).

Стоимость: измеряем, а не верим

Тот же счётчик, четыре реализации, 40 миллионов инкрементов. Порядок величин на восьмиядерной машине (у вас будут другие числа — важны соотношения):

Реализация Время Что происходит
один поток, обычный u64 ~40 мс инкремент в регистре, никакой синхронизации
8 потоков, Arc<Mutex<u64>> ~5000 мс каждый инкремент — CAS плюс возможный сон в ядре; в 100+ раз хуже однопоточного
8 потоков, AtomicU64::fetch_add(Relaxed) ~1500 мс без системных вызовов, но кэш-линия ходит между ядрами
8 потоков, локальный счётчик + один fetch_add в конце ~10 мс синхронизация 8 раз вместо 40 миллионов

Выводы, которые стоит выписать на стену:

  1. Конкурентность не ускоряет автоматически. Разделяемая изменяемая ячейка — это точка сериализации; закон Амдала считает именно её («Производительность конкурентности»).
  2. Лучшая синхронизация — её отсутствие. Схема «работаем в локальных данных, сливаем результат один раз» бьёт любую оптимизацию замков. Это же принцип sharding: разбить состояние на N частей со своим замком (dashmap устроен так).
  3. Мьютекс дорог не сам по себе, а под конкуренцией. Незанятый lock() — это один атомарный CAS, десятки наносекунд. Занятый — futex, планировщик, переключение контекста, микросекунды.
  4. Прежде чем оптимизировать — профилируйте: perf покажет время в futex_wait, а perf c2c — межъядерный трафик («Профилирование CPU»).

Атомики и порядки памяти

Атомарные типы (AtomicUsize, AtomicU64, AtomicBool, AtomicPtr) — это операции, которые процессор выполняет неделимо, и они же — единственный легальный способ иметь изменяемое состояние без замка. У каждой операции есть параметр Ordering, и он не про скорость, а про видимость соседних, неатомарных записей.

use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::thread;

fn main() {
    let hits = Arc::new(AtomicU64::new(0));

    let mut hs = Vec::new();
    for _ in 0..8 {
        let hits = Arc::clone(&hits);
        hs.push(thread::spawn(move || {
            for _ in 0..1000 {
                // Relaxed достаточно: нам нужна только атомарность самого счётчика,
                // никакие другие данные к нему не привязаны
                hits.fetch_add(1, Ordering::Relaxed);
            }
        }));
    }
    for h in hs { h.join().unwrap(); }
    println!("{}", hits.load(Ordering::Relaxed));   // 8000
}

А вот случай, где Relaxed — ошибка. Один поток готовит данные и поднимает флаг, другой ждёт флаг и читает данные:

use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use std::thread;

struct Shared { value: AtomicU64, ready: AtomicBool }

fn main() {
    let s = Arc::new(Shared { value: AtomicU64::new(0), ready: AtomicBool::new(false) });

    let writer = { let s = Arc::clone(&s); thread::spawn(move || {
        s.value.store(42, Ordering::Relaxed);
        s.ready.store(true, Ordering::Release);   // ВСЁ, записанное выше, становится видимым
    })};

    let reader = { let s = Arc::clone(&s); thread::spawn(move || {
        while !s.ready.load(Ordering::Acquire) {  // ...тому, кто прочтёт true через Acquire
            std::hint::spin_loop();               // подсказка процессору: это спин
        }
        assert_eq!(s.value.load(Ordering::Relaxed), 42);   // гарантированно 42
    })};

    writer.join().unwrap();
    reader.join().unwrap();
}

Замените Release/Acquire на Relaxed — и assert может упасть, увидев 0: процессор и компилятор вправе переупорядочить две независимые записи. Обратите внимание на важное: UB здесь не возникает — атомик остаётся атомиком. Вы просто получаете неверный результат, и Rust вас от этого не спасает. Это ровно тот случай, где «бесстрашность» заканчивается.

Практическая дисциплина:

  • SeqCst по умолчанию. Он самый строгий и самый понятный. На x86-64 разница с Acquire/Release для загрузок и сохранений почти нулевая, платите вы в основном на ARM.
  • Relaxed — только для независимых счётчиков и статистики, где ни от чего не зависит порядок.
  • Release в сохранении, Acquire в загрузке — всегда парой. Одиночный Release бессмыслен: он публикует данные для того, кто читает Acquire.
  • compare_exchange(current, new, success, failure) — базовый строительный блок для «прочитал, посчитал, записал, если не изменилось»:
use std::sync::atomic::{AtomicU64, Ordering};

// Атомарно записать максимум: классический CAS-цикл
fn store_max(cell: &AtomicU64, candidate: u64) {
    let mut cur = cell.load(Ordering::Relaxed);
    while candidate > cur {
        match cell.compare_exchange_weak(cur, candidate, Ordering::Release, Ordering::Relaxed) {
            Ok(_) => return,
            Err(actual) => cur = actual,   // кто-то опередил — пробуем снова с новым значением
        }
    }
}

Полноценные lock-free структуры (очереди, списки, hazard pointers) в безопасном Rust не пишутся: там нужен unsafe и решение проблемы освобождения памяти — этим занимаются crossbeam-epoch и crossbeam::queue. Честный совет: не пишите свои lock-free структуры. Возьмите готовую (crossbeam, dashmap, flume, arc-swap) или используйте мьютекс на шардированном состоянии — почти всегда этого достаточно.

Ложное разделение: два счётчика в одной кэш-линии

Отдельная засада, которую не поймает никакой компилятор, — ложное разделение. Два атомика в соседних полях структуры попадают в одну 64-байтную кэш-линию, и хотя логически они независимы, физически ядра дерутся за линию. Лечится crossbeam_utils::CachePadded или, лучше, полным отказом от общего счётчика в горячем цикле. Тема разбирается в «Кэш и локальность».

Глобальное состояние без static mut

static mut требует unsafe при каждом обращении, а в редакции 2024 ещё и отдельно ругается линтом static_mut_refs. Правильные инструменты:

use std::collections::HashMap;
use std::sync::{LazyLock, Mutex, OnceLock};

// 1. Конфигурация, которая заполняется один раз при старте и потом только читается
static CONFIG: OnceLock<String> = OnceLock::new();

fn config() -> &'static str {
    CONFIG.get_or_init(|| std::env::var("APP_MODE").unwrap_or_else(|_| "dev".into()))
}

// 2. Ленивая инициализация с вычислением (замена крейта lazy_static / once_cell)
static TABLE: LazyLock<HashMap<&str, u32>> =
    LazyLock::new(|| HashMap::from([("gzip", 1), ("br", 2)]));

// 3. Изменяемое глобальное состояние: только через синхронизацию
static COUNTERS: LazyLock<Mutex<HashMap<String, u64>>> =
    LazyLock::new(|| Mutex::new(HashMap::new()));

// 4. Состояние на поток: без синхронизации вообще, потому что нечего разделять
thread_local! {
    static SCRATCH: std::cell::RefCell<Vec<u8>> = std::cell::RefCell::new(Vec::new());
}

fn reuse_buffer() -> usize {
    SCRATCH.with_borrow_mut(|buf| { buf.clear(); buf.extend_from_slice(b"data"); buf.len() })
}

thread_local! — недооценённый инструмент: переиспользуемый буфер, генератор случайных чисел, арена парсера. Ноль синхронизации, потому что данные не разделяются. Оговорка: деструкторы thread_local не запускаются для главного потока при завершении процесса и вообще не самая надёжная точка для важной очистки.

Данные, а не задачи: rayon

Если задача формулируется как «примени это ко всем элементам», ручные потоки не нужны. rayon даёт параллелизм по данным через work-stealing пул: меняете iter() на par_iter(), и работа раскладывается по ядрам.

use rayon::prelude::*;

fn main() {
    let data: Vec<u64> = (1..=20_000_000).collect();

    // Последовательно: sum = data.iter().filter(...).map(...).sum();
    let sum: u64 = data
        .par_iter()                          // требует: элементы Send + Sync
        .filter(|&&x| x % 3 == 0)
        .map(|&x| x * x % 1_000_003)
        .sum();                              // reduce делается по дереву, а не в одну ячейку

    println!("{sum}");
}

Почему это не превращается в кашу: par_iter требует Send + Sync, а замыкания получают только &T, поэтому написать гонку через rayon в безопасном коде невозможно — компилятор не даст. Что нужно помнить:

  • sum, reduce, fold работают деревом, поэтому операция должна быть ассоциативной. Для чисел с плавающей точкой это означает, что результат может отличаться от последовательного в последних битах — воспроизводимость требует детерминированного порядка (par_chunks с ручной сборкой).
  • Мелкая работа не окупается. На теле цикла в несколько наносекунд накладные расходы разбиения съедят всё; rayon адаптивен, но не всемогущ. Меряйте.
  • rayon и блокирующий ввод-вывод несовместимы. Пул размером по числу ядер, забитый ожиданием сети, — это простой. Для I/O — асинхронность.
  • rayon::join(a, b) — для рекурсивных «разделяй и властвуй» (сортировки, обходы деревьев); rayon::scope — для произвольного fork-join.

Типичные грабли

1. Забыли drop(tx) — программа не завершается. Разобрано выше. Симптом: последний recv() висит вечно. Проверка: сколько живых Sender существует в момент, когда ожидаете конец?

2. move в цикле — tx перемещён на первой итерации.

error[E0382]: use of moved value: `tx`
   --> src/main.rs:9:23
    |
9   |         thread::spawn(move || { tx.send(id).unwrap(); });
    |                       ^^^^^^^   -- use occurs due to use in closure
    |                       |
    |                       value moved into closure here, in previous iteration of loop
    |
help: consider cloning the value if the performance cost is acceptable

Фраза in previous iteration of loop — точное указание причины. Лечение: let tx = tx.clone(); внутри цикла до spawn. Тот же приём для Arc: let x = Arc::clone(&x);. Идиома с одноимённой переменной (shadowing) здесь стандартна и читается лучше, чем tx2.

3. Guard живёт дольше критической секции. Самая дорогая ошибка производительности:

// ПЛОХО: тяжёлая работа под замком — все остальные потоки стоят
let mut cache = cache.lock().unwrap();
let value = expensive_computation(&input);   // 50 мс с запертым мьютексом
cache.insert(input, value);

// ХОРОШО: считаем вне замка, под замком только вставка
let value = expensive_computation(&input);
cache.lock().unwrap().insert(input, value);  // временный guard умирает в конце инструкции

Особый случай — .await при живом guard: MutexGuard не Send, поэтому задача перестаёт быть Send и компилятор ругается длинным сообщением. Разбор — в «Асинхронном Rust».

4. Двойной захват одного мьютекса в одном выражении. std::sync::Mutex не реентрантный: повторный lock() из того же потока — это вечное ожидание себя.

// ЗАВИСНЕТ: временный guard от первого lock() жив до конца всего match
match data.lock().unwrap().len() {
    0 => println!("пусто"),
    n => data.lock().unwrap().push(n),   // второй lock в том же потоке
}
// Правильно: let n = data.lock().unwrap().len(); а потом match n { ... }

То же самое случается, когда метод под замком вызывает другой публичный метод, который тоже берёт замок. Правило: приватные функции принимают уже полученный &mut T, замок берут только публичные точки входа.

5. Взаимоблокировка на двух замках. Канонический пример — перевод денег:

use std::sync::Mutex;

struct Account { id: u32, balance: Mutex<i64> }

// ПЛОХО: transfer(a, b) и transfer(b, a) параллельно — гарантированный взаимный клинч
fn transfer_bad(from: &Account, to: &Account, amount: i64) {
    let mut f = from.balance.lock().unwrap();
    let mut t = to.balance.lock().unwrap();   // ждём замок, который держит второй поток
    *f -= amount;
    *t += amount;
}

// ХОРОШО: единый глобальный порядок захвата — по возрастанию id
fn transfer(from: &Account, to: &Account, amount: i64) {
    assert_ne!(from.id, to.id, "перевод самому себе обрабатывается отдельно");
    let (first, second) = if from.id < to.id { (from, to) } else { (to, from) };
    let mut g1 = first.balance.lock().unwrap();
    let mut g2 = second.balance.lock().unwrap();
    let (f, t) = if first.id == from.id { (&mut *g1, &mut *g2) } else { (&mut *g2, &mut *g1) };
    *f -= amount;
    *t += amount;
}

Четыре рабочих способа не иметь взаимоблокировок: (а) держать один замок за раз; (б) глобальный порядок захвата, задокументированный в коде; (в) try_lock с откатом и повтором; (г) не разделять состояние вообще — актор с каналом владеет данными единолично.

6. Паника в потоке теряется. join() вернёт Err, но если вы не вызываете join() или пишете let _ = h.join();, поток умрёт тихо. В сервисе это выглядит как «воркер перестал брать задачи». В проде ставьте std::panic::set_hook, который пишет в лог со thread::current().name(), и следите за живостью воркеров («Прод на Rust»). И помните: panic = "abort" в профиле release превращает панику любого потока в смерть процесса — иногда это именно то, что нужно.

7. Отравление мьютекса каскадом. Один поток паниковал под замком — все последующие lock().unwrap() паникуют тоже, и падает весь сервис. Варианты: parking_lot (нет отравления), либо явное решение о судьбе данных:

let mut guard = mtx.lock().unwrap_or_else(|e| e.into_inner());   // осознанно берём как есть

8. Поток на задачу. «Пришёл запрос — создали поток» держится до сотен RPS, потом падает под весом стеков и переключений контекста. Правильно — пул фиксированного размера (rayon, threadpool) или асинхронный рантайм.

9. Активное ожидание. while !flag.load(Relaxed) {} съедает ядро целиком и мешает тому потоку, который должен поднять флаг. Спин оправдан только для ожидания в наносекунды и только с std::hint::spin_loop(). Всё остальное — канал, Condvar, thread::park/unpark. Из той же серии: логирование под замком и Instant::now() на каждой итерации — частая причина, почему параллельная версия оказалась медленнее последовательной.

Как выбрать примитив

Отладка и тестирование конкурентного кода

Конкурентные баги не ловятся чтением кода, поэтому инструменты нужны заранее.

  • Висим — снимите стеки. gdb -p PID, затем thread apply all bt. Имена потоков (Builder::name) делают вывод читаемым. rust-gdb/rust-lldb из состава rustup умеют печатать Vec и String осмысленно («Отладка»).
  • perf покажет, где время: много futex_wait — конкуренция за замок; много cache-misses при росте потоков — ложное разделение; perf c2c находит конкретные кэш-линии.
  • parking_lot с фичей deadlock_detection периодически проверяет граф ожиданий и печатает цикл — дешёвая страховка на стенде.
  • loom перебирает возможные перестановки операций в модели памяти и находит баги, которые «никогда не воспроизводятся». Применим к крошечным моделям (два потока, три операции) и нужен тем, кто пишет свои примитивы. См. «Тестирование».
  • Miri умеет исполнять многопоточный код и ловит гонки данных и UB в unsafe-блоках: cargo +nightly miri test. Медленно, но незаменимо.
  • ThreadSanitizerRUSTFLAGS="-Zsanitizer=thread" cargo +nightly test — ловит гонки в FFI и unsafe на реальных прогонах.
  • Тесты с таймаутом. Конкурентный тест, который висит, ломает CI хуже упавшего. Оборачивайте ожидание в recv_timeout, а не recv.

Цена и границы: честно

Где Rust реально выигрывает. Многопоточный код, который живёт годами и меняется многими руками. Именно здесь ценность гарантии максимальна: новый человек не может случайно добавить гонку, а рефакторинг «вынести это в поток» либо компилируется, либо честно объясняет, почему нет. Второй выигрыш — отсутствие GC-пауз: в C# и Java конкурентный сборщик добавляет хвост латентности, которого в Rust нет («JVM и память»).

Где больно. Структуры данных с нетривиальным разделением (lock-free очереди, кэши с обратными ссылками, наблюдатели) требуют unsafe, Pin или арен. Отладка отравления и «почему это не Send» отнимает время, которого в Go или Java вы бы не потратили — там бы просто получили редкий баг вместо ошибки компиляции; это обмен, а не подарок. Компиляция многопоточного кода с rayon и дженериками — минуты, и это ощутимо давит на цикл обратной связи.

Где Rust не тот инструмент. Массовая изоляция отказов и горячая перезагрузка кода — это BEAM, и конкурировать с супервизорами Elixir на его поле не стоит («Конкурентность в Elixir»). Сотни тысяч одновременных дешёвых соединений при простом коде — Go тут даёт близкий результат гораздо быстрее в разработке. Скрипт, который нужно распараллелить один раз, — это xargs -P или пул процессов в Python.

И главное. Rust не делает конкурентность простой. Он делает её отлаживаемой: убирает невоспроизводимый класс ошибок и оставляет вам архитектурные решения — что разделять, в каком порядке брать замки, где ставить границы очередей. Эти решения по-прежнему трудные, и никакой компилятор их не примет за вас.

Связь с соседними треками

Мини-итог

  • «Бесстрашная конкурентность» — точное утверждение: гонки данных невыразимы в безопасном Rust, и это следствие правила заимствования, а не отдельного механизма.
  • Компилятор ловит гонки данных, use-after-free и «забыл замок». Он не ловит взаимоблокировки, логические состояния гонки, неверные Ordering, голодание и неограниченные очереди.
  • Send — «можно переместить в поток». Sync — «можно раздать &T многим потокам», то есть &T: Send. Оба выводятся структурно; Rc не Send, RefCell не Sync, MutexGuard не Send.
  • Читайте E0277 снизу вверх: граница spawn → замыкание → захваченное поле → виновный тип.
  • thread::spawn требует 'static; thread::scope снимает требование, гарантируя завершение потоков до выхода из блока.
  • Канал вместо разделения, когда возможно: send передаёт владение. drop(tx) обязателен, ограниченный канал — по умолчанию.
  • Arc<Mutex<T>> — это «разделяемое владение» плюс «синхронизация», два ортогональных типа. Главное свойство: Mutex владеет данными, поэтому доступ без замка невозможен по построению.
  • Критическая секция должна быть минимальной. Тяжёлые вычисления и .await — вне замка.
  • Атомики — для счётчиков и флагов. SeqCst по умолчанию, Relaxed осознанно, Release/Acquire только парой. Свои lock-free структуры не писать.
  • Лучшая синхронизация — отсутствие синхронизации: локальные накопители, шардирование, thread_local, rayon по данным.
  • Измеряйте. Восемь потоков на одном мьютексе легко работают в сто раз медленнее одного потока без него.

Источники

Что дальше

Асинхронный Rust: future, tokio, async/await и его подводные камни — в этой главе мы дважды упирались в одно ограничение: поток ОС слишком дорог, чтобы выделять его на ожидание. Тысяча соединений — не тысяча потоков. Следующая глава показывает второй ответ Rust на конкурентность: Future как конечный автомат, который компилятор собирает за вас, tokio как планировщик этих автоматов, и цена этого решения — от «почему мой MutexGuard не проходит через .await» до расщепления экосистемы на синхронную и асинхронную половины.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов