← Назад к списку тем

14. Конкурентность: продвинутый уровень

Атомики и memory ordering, lock-free структуры, scoped threads, rayon, work-stealing, false sharing.

Атомики — операции без блокировок

std::sync::atomic предоставляет типы, чьи операции транслируются напрямую в атомарные инструкции процессора (например, lock xadd на x86), минуя ОС-планировщик и очередь ожидания. Для простых счётчиков и флагов это на порядок быстрее Mutex.

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

let counter = Arc::new(AtomicUsize::new(0));
let mut handles = vec![];

for _ in 0..8 {
    let counter = Arc::clone(&counter);
    handles.push(thread::spawn(move || {
        for _ in 0..100_000 {
            counter.fetch_add(1, Ordering::Relaxed);
        }
    }));
}
for h in handles { h.join().unwrap(); }
println!("{}", counter.load(Ordering::Relaxed)); // 800000

Memory ordering — что на самом деле означает Ordering

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

// Relaxed — только атомарность, никаких гарантий порядка с другими операциями.
// Годится для независимых счётчиков (метрики, статистика).
counter.fetch_add(1, Ordering::Relaxed);

// Acquire/Release — пара для передачи данных между потоками через флаг.
// Release: все записи ДО этой операции видны потоку, который сделает Acquire ПОСЛЕ.
data.store(42, Ordering::Relaxed);
ready.store(true, Ordering::Release);      // публикуем данные

// В другом потоке:
if ready.load(Ordering::Acquire) {          // синхронизируется с Release выше
    println!("{}", data.load(Ordering::Relaxed)); // гарантированно увидит 42
}

// SeqCst — самый строгий: глобальный тотальный порядок для ВСЕХ SeqCst операций.
// Проще рассуждать, но дороже — компилятор не может некоторые оптимизации применить.
flag.store(true, Ordering::SeqCst);
OrderingГарантияСтоимость
Relaxedтолько атомарность операцииминимальная
Acquire / Releaseоднонаправленный барьер (load/store)средняя
AcqRelAcquire + Release вместе (для RMW-операций)средняя
SeqCstглобальный тотальный порядокмаксимальная
⚠️ Не гадайте с ordering: если не уверены, какой memory ordering нужен — начните с SeqCst для корректности, а затем ослабляйте до Acquire/Release/Relaxed только с профилированием и пониманием паттерна доступа. Ошибка в ordering — undefined behavior, которое может годами не проявляться на x86 (сильная модель памяти) и сразу же ломаться на ARM (слабая модель).

Compare-and-swap и lock-free структуры

CAS (compare-and-swap) — строительный блок lock-free алгоритмов: атомарно сравнивает текущее значение с ожидаемым и, если совпало, заменяет на новое. На этом строятся lock-free стеки, очереди, счётчики без блокировок вообще.

use std::sync::atomic::{AtomicI64, Ordering};

fn update_max(current_max: &AtomicI64, candidate: i64) {
    let mut old = current_max.load(Ordering::Relaxed);
    loop {
        if candidate <= old { return; }
        // compare_exchange: если значение всё ещё old — заменить на candidate
        match current_max.compare_exchange(
            old, candidate, Ordering::Relaxed, Ordering::Relaxed
        ) {
            Ok(_) => return,               // успех
            Err(actual) => old = actual, // кто-то опередил — повторяем с новым old
        }
    }
}

Этот паттерн — "retry loop" — типичен для lock-free кода: вместо блокировки поток крутится в цикле, пока его CAS не пройдёт успешно. При высокой contention это может быть хуже мьютекса (много впустую потраченных CAS), поэтому lock-free не значит "всегда быстрее".

🚫 Не пишите lock-free структуры с нуля в продакшене: корректные lock-free алгоритмы (особенно очереди, требующие решения ABA-проблемы) крайне тяжело верифицировать. Используйте проверенные крейты — crossbeam (SegQueue, ArrayQueue, epoch-based reclamation) вместо самописного кода.

Scoped threads

До Rust 1.63 заимствовать данные из внешнего стека в поток было невозможно без Arc — компилятор не мог доказать, что поток завершится раньше, чем данные выйдут из области видимости. std::thread::scope решает это: гарантирует (через API, а не unsafe), что все потоки внутри scope join'атся перед выходом.

use std::thread;

let data = vec![1, 2, 3];

thread::scope(|s| {
    s.spawn(|| {
        println!("первая половина: {:?}", &data[..1]);
    });
    s.spawn(|| {
        println!("вторая половина: {:?}", &data[1..]);
    });
});  // scope блокируется здесь, пока все spawn'ы не завершатся

println!("data всё ещё доступна: {:?}", data); // не перемещалась!

Никакого Arc, никакого move — обычные заимствования &data работают, потому что компилятор видит границу scope как гарантированную точку join.

rayon — параллелизм данных и work-stealing

rayon — де-факто стандартная библиотека для параллелизма данных (data parallelism). Превращает последовательные итераторы в параллельные заменой одного вызова, используя пул потоков с work-stealing планировщиком: свободный поток "ворует" задачи из очереди занятого, минимизируя простой.

use rayon::prelude::*;

let numbers: Vec<u64> = (0..10_000_000).collect();

// Было (последовательно):
let sum: u64 = numbers.iter().map(|&x| x * x).sum();

// Стало (параллельно) — просто par_iter() вместо iter()
let sum: u64 = numbers.par_iter().map(|&x| x * x).sum();

// par_sort — параллельная сортировка "на месте"
let mut v = numbers.clone();
v.par_sort_unstable();

// join — рекурсивное разделение задач (divide and conquer)
fn quicksort<T: Ord + Send>(v: &mut [T]) {
    if v.len() <= 1 { return; }
    let mid = partition(v);
    let (left, right) = v.split_at_mut(mid);
    rayon::join(|| quicksort(left), || quicksort(right));
}

rayon безопасен благодаря тем же Send/Sync гарантиям на уровне типов: замыкание не скомпилируется, если захватывает что-то, небезопасное для параллельного доступа.

False sharing — невидимый убийца производительности

Кэш-линия процессора обычно 64 байта. Если два атомика/поля, используемых разными потоками, физически попадают в одну кэш-линию, запись в одно поле инвалидирует кэш другого потока — даже если логически данные независимы. Это называется false sharing: код корректен, но производительность падает в разы.

// Плохо: два атомика в соседних полях — вероятно, одна кэш-линия
struct Counters {
    reads: AtomicUsize,   // поток A постоянно пишет сюда
    writes: AtomicUsize,  // поток B постоянно пишет сюда — false sharing!
}

// Лучше: выравнивание по границе кэш-линии разносит поля в разные линии
#[repr(align(64))]
struct PaddedCounter {
    value: AtomicUsize,
}

struct Counters {
    reads: PaddedCounter,
    writes: PaddedCounter,   // теперь гарантированно в разных кэш-линиях
}

// crossbeam::utils::CachePadded<T> делает то же самое из коробки
use crossbeam::utils::CachePadded;
struct Counters2 {
    reads: CachePadded<AtomicUsize>,
    writes: CachePadded<AtomicUsize>,
}
🔑 Ключевое: false sharing не ловится ни компилятором, ни обычным профилировщиком по CPU-времени — нужны счётчики кэш-промахов (perf stat, cache-misses). Подозревайте его, если параллельный код с независимыми данными масштабируется хуже линейного с ростом числа потоков.

Выбор инструмента: сравнение подходов

ПодходКогда использоватьСтоимость
Mutex<T>сложные инварианты, редкие обновлениясредняя (syscall при contention)
Atomic*простые счётчики, флаги, версиинизкая
crossbeam lock-free структурывысокая contention, готовый алгоритм подходитнизкая, но сложнее в отладке
thread::scopeпараллельная обработка данных без 'staticминимальная (обычные потоки ОС)
rayonCPU-bound задачи, легко разбиваемые на подзадачинизкая, автоматический work-stealing

Golden rule: начинайте с самого простого решения (Mutex), переходите к атомикам и lock-free только когда профилирование показало contention на блокировке как узкое место — преждевременная lock-free оптимизация усложняет код без гарантированной выгоды.