Бесстрашная конкурентность: потоки, каналы, 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 в регистре и вынести его из цикла — тогда результат станет вообще произвольным («Неопределённое поведение»).
Дадим точное определение, потому что дальше всё держится на нём. Гонка данных — это одновременное выполнение трёх условий:
- два или более потока обращаются к одной ячейке памяти;
- хотя бы одно обращение — запись;
- обращения не упорядочены отношением 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 уничтожены, — и отдаёт значение назад, чтобы вы могли его залогировать или переложить в другую очередь. Ни один байт не теряется молча.
процессор не тратится P2->>Q: send(Job 2) Q-->>C: Ok(Job 2) — владение переходит consumer Note over C: Job принадлежит только consumer:
синхронизация больше не нужна Q-->>C: Ok(Job 1) P1->>Q: drop(Sender) P2->>Q: drop(Sender) Note over Q: живых отправителей ноль Q-->>C: Err(RecvError) — цикл for по rx завершается штатно
Ограниченный канал и обратное давление
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 это два разных типа, и понимание, зачем каждый, — половина дела.
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 миллионов |
Выводы, которые стоит выписать на стену:
- Конкурентность не ускоряет автоматически. Разделяемая изменяемая ячейка — это точка сериализации; закон Амдала считает именно её («Производительность конкурентности»).
- Лучшая синхронизация — её отсутствие. Схема «работаем в локальных данных, сливаем результат один раз» бьёт любую оптимизацию замков. Это же принцип sharding: разбить состояние на N частей со своим замком (
dashmapустроен так). - Мьютекс дорог не сам по себе, а под конкуренцией. Незанятый
lock()— это один атомарный CAS, десятки наносекунд. Занятый — futex, планировщик, переключение контекста, микросекунды. - Прежде чем оптимизировать — профилируйте:
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. Медленно, но незаменимо. - ThreadSanitizer —
RUSTFLAGS="-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 не делает конкурентность простой. Он делает её отлаживаемой: убирает невоспроизводимый класс ошибок и оставляет вам архитектурные решения — что разделять, в каком порядке брать замки, где ставить границы очередей. Эти решения по-прежнему трудные, и никакой компилятор их не примет за вас.
Связь с соседними треками
- Как это устроено на уровне ОС — «Потоки и синхронизация» (те же грабли на C, где их не ловит никто) и «Процессы и планирование»: futex, переключение контекста, что стоит парковка потока.
- Почему параллельное не значит быстрое — «Производительность конкурентности» и «Кэш и локальность»: закон Амдала, ложное разделение, стоимость когерентности.
- Тот же вопрос в масштабе сети — «Координация»: распределённые замки — это те же проблемы плюс отказы и задержки, и именно поэтому их стараются не использовать.
- Другие ответы на ту же задачу — «Конкурентность в Go» (горутины и CSP с рантаймом) и «Конкурентность в Elixir» (изолированные процессы и обмен копиями).
Мини-итог
- «Бесстрашная конкурентность» — точное утверждение: гонки данных невыразимы в безопасном 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по данным. - Измеряйте. Восемь потоков на одном мьютексе легко работают в сто раз медленнее одного потока без него.
Источники
- The Rust Programming Language, глава 16 «Fearless Concurrency» — канонический вход: потоки, каналы,
Mutex,Send/Sync. - Mara Bos, «Rust Atomics and Locks» — бесплатная онлайн-книга, лучший существующий текст про атомики, модель памяти и реализацию примитивов; автор — руководитель библиотечной команды Rust.
- The Rustonomicon: Send and Sync — что именно вы обещаете, когда пишете
unsafe impl. - Документация
std::sync,std::threadиstd::sync::atomic— читать вместе с разделами «Poisoning» и «Memory Ordering». - Aaron Turon, «Fearless Concurrency with Rust» (2015) — исходная статья, из которой пошёл термин; полезна как исторический контекст.
crossbeam,rayon,parking_lot,dashmap— четыре крейта, покрывающие 90% практических нужд.loomи Miri — проверка конкурентного кода перебором перестановок и интерпретацией.- Hans Boehm, «Threads Cannot be Implemented as a Library» — почему модель памяти обязана быть частью языка, а не библиотеки.
- Jeff Preshing, «An Introduction to Lock-Free Programming» и «Acquire and Release Semantics» — лучшее объяснение порядков памяти вне контекста Rust.
- Jon Gjengset, «Rust for Rustaceans» — главы про конкурентность и
unsafe-абстракции; его видеоразборы с живым написанием конкурентных структур особенно полезны. - Rust Atomics:
std::sync::mpscпосле переписывания на crossbeam — как менялась реализация каналов в стандартной библиотеке.
Что дальше
Асинхронный Rust: future, tokio, async/await и его подводные камни — в этой главе мы дважды упирались в одно ограничение: поток ОС слишком дорог, чтобы выделять его на ожидание. Тысяча соединений — не тысяча потоков. Следующая глава показывает второй ответ Rust на конкурентность: Future как конечный автомат, который компилятор собирает за вас, tokio как планировщик этих автоматов, и цена этого решения — от «почему мой MutexGuard не проходит через .await» до расщепления экосистемы на синхронную и асинхронную половины.