Съдържание на курса
Урок 7 — Конкурентност и async
От Dorian Chávez · основател на Hábil и архитект на интеграции ·
Време: 2 × 45 мин.
Какво изграждаш: конкурентният revisor: да проверява всичко едновременно
Какво научаваш: нишки на операционната система, Arc и Mutex, канали, async/await с tokio, честното сравнение с goroutine-ите на Go
След края ще можеш да
- Стартираш нишки на операционната система с
thread::spawn, да им прехвърляш правилната собственост и да взимаш резултата им сjoin. - Обясниш защо една данна, споделена между нишки, се нуждае от
Arc, кога се нуждае допълнително отMutexи какво защитаваMutexGuard. - Използваш канал от стандартната библиотека, за да предаваш резултати, без да споделяш променлива колекция.
- Четеш грешката
E0277, когато една стойност не отговаря наSend, и да разпознаваш, че компилаторът защитава граница между нишки. - Обясниш разликата между конкурентност и паралелизъм и между нишка на операционната система, async задача и goroutine на Go.
- Следваш в
revisorпътя на едно async запитване, от#[tokio::main]доjoin_all, семафора и състоянието, което остава свързано с всяка услуга.
Защо, преди как
Досега revisor може да получи списък със услуги, да пита една и да класифицира отговора. Ако пита десет услуги последователно и всяка се бави по една секунда да отговори, докладът отнема приблизително десет секунди. Няма значение, че компютърът има няколко ядра, нито че програмата е бърза: през почти цялото това време процесорът не изчислява нищо; чака отговор от мрежата.
Чакането на отговор от мрежата е различно от изчисляването на голяма сума. При интензивно изчисление повече ядра могат да позволят наистина паралелна работа. При едно HTTP запитване, напротив, по-голямата част от времето е извън процеса: операционната система чака пакети, отдалеченият сървър решава какво да прави и мрежата пренася отговора. Междувременно програмата би могла да започне други запитвания. Това е конкурентност: организиране на няколко задачи, които напредват в редуващи се периоди. Може да се превърне в паралелизъм, ако няколко задачи изпълняват код едновременно на различни ядра, но не са синоними.
Практическата цел на този урок е да промени общото време. Ако пет услуги се бавят близо по една секунда всяка и ги питаш една по една, докладът отнема около пет секунди. Ако започнеш петте запитвания и чакаш отговорите им едновременно, времето се доближава до това на най-бавната услуга, не до сбора на всички. Това не кара отдалечените услуги да отговарят по-бързо; избягва се губенето на времето, което програмата е прекарвала в чакане на първата, преди да започне втората.
Цената е, че няколко части на програмата могат да са живи едновременно. Появяват се въпроси, които последователният код нямаше: кой е собственикът на данната? кога приключва една задача? какъв ред имат резултатите? могат ли две задачи да променят една и съща стойност? какво става, ако една се провали? как не позволяваш да се отворят хиляди едновременни връзки? Rust не отговаря на тези въпроси, като скрива споделената памет. Прави така, че ownership, заемания и trait-ове като Send и Sync да продължават да имат значение, когато има няколко нишки или задачи.
Това се свързва пряко с урок 2. Ownership изглеждаше като локално правило: една стойност има собственик и заеманията трябва да уважават живота ѝ. При конкурентността тези правила се превръщат в гаранция между задачите. Една нишка не може да запази референция към променлива на main, която може вече да е изчезнала. Не може и да получи тип, за който Rust знае, че не е безопасно да бъде преместен в друга нишка. Това, което в началото се усещаше като триене с компилатора, тук се превръща в преграда срещу висящи референции и състезания за данни (data race) в безопасен код.
Rust предлага два основни инструмента за проблема. Стандартната библиотека носи нишки на операционната система, mutex-и, атомарни броячи на референции и канали. Подходящи са за работа на процесора, малки програми или интеграция с блокиращи API-та. За много мрежови операции, които прекарват време в чакане, revisor използва async/await и Tokio. Async не означава „по-бързо по дефиниция“: означава, че умерен брой нишки могат да придвижват много операции, които чакат вход-изход, без да резервират една блокирана нишка за всяка.
Go взима друго решение. Една goroutine се стартира с go f() и runtime-ът е интегриран в езика и в дистрибуцията. Rust те принуждава да различаваш нишка от async задача и, за async, да избереш runtime. Това изисква повече речник и повече решения, но позволява типът на данните и границата на собствеността да са явни. Никой от двата подхода не премахва нуждата да се проектират граници, да се обработват грешки и да се измерва. Полезното сравнение не е кой език „печели“, а каква цена плаща всеки и каква гаранция предлага в замяна.
Първо прочети глави 16 и 17 на The Rust Book. После направи упражненията на Rustlings за нишки и канали, преди да адаптираш revisor. Редът има значение: async става много по-малко загадъчен, когато вече разбираш какво означава един closure да бъде преместен в друга нишка, какво означава Send и защо споделянето на променливост изисква видима синхронизация.
Понятията
Нишки на операционната система, move и join
std::thread::spawn иска closure и стартира нишка на операционната система, за да го изпълни. Стойността, която връща, е JoinHandle<T>: конкретно обещание, че нишката може да завърши със стойност от тип T. Извикването на join() чака тази нишка да приключи и връща Result<T, Box<dyn Any + Send>>; Err представя, че нишката е направила panic!.
Ключовата дума move е важна. Един closure без move може да се опита да улови референция към променлива от външния контекст. Една нишка може да продължи да се изпълнява, след като този контекст е приключил, така че Rust не позволява да се предаде на нишката референция, която може да престане да е валидна. move кара closure-а да улавя по стойност. На фигурата всяко Servicio престава да принадлежи на вектора и започва да принадлежи на closure-а на собствената си нишка.
Фиг. 7.1 | Една нишка на услуга.
// fig07_01.rs
use std::thread;
struct Servicio {
nombre: String,
}
type Estado = String; // en el curso es el enum de la lección 3; aquí basta un texto
fn revisar(s: &Servicio) -> Estado {
format!("{}: OK", s.nombre)
}
fn main() {
let servicios = vec![
Servicio { nombre: "catalogo".to_string() },
Servicio { nombre: "pagos".to_string() },
Servicio { nombre: "reportes".to_string() },
];
let handles: Vec<_> = servicios.into_iter().map(|s| {
thread::spawn(move || revisar(&s)) // `move` entrega la propiedad al hilo
}).collect();
for h in handles {
let estado = h.join().unwrap(); // espera y recoge el resultado
println!("{estado}");
}
}
$ rustc --edition 2024 fig07_01.rs && ./fig07_01
catalogo: OK
pagos: OK
reportes: OK
Редът на println! е детерминиран, дори нишките да завършват в друг ред. JoinHandle-овете се съхраняват в същия ред, който произвежда servicios.into_iter(), а вторият for извиква join() в този ред. Ако нишката на pagos приключи първа, стойността ѝ е готова, но програмата първо чака и отпечатва тази на catalogo. Това различие има значение: конкурентността не задължава изходът да е недетерминиран. Можеш да проектираш граница, където наблюдаваният ред остава стабилен.
join е достатъчен, когато всяка задача има резултат и броят на задачите е малък и известен. Не ти трябва канал само за да вземеш една стойност; handle-ът вече я предава. В Go обикновено би комбинирал една goroutine с WaitGroup, за да чакаш, и канал или защитена колекция, за да вземеш резултатите. Rust прави резултата част от handle-а, макар това да не премахва полезността на каналите за прогресивна комуникация.
Една нишка на операционната система не е безплатна. Има ресурси на операционната система, стек и цена за планиране, по-голяма от тази на една async задача. Няма универсална цифра за памет на нишка: зависи от операционната система, архитектурата и конфигурацията. Правилото за дизайн е по-полезно от фиксирано число: не отваряй нишка на операционната система за всяка връзка, ако основната работа се състои в чакане на мрежа. За няколко задачи на процесора или блокиращ API една нишка може да е точно подходящото. За много HTTP услуги revisor използва async.
revisor не създава нишка на услуга. Мрежовата му работа се изразява като async future-и; runtime-ът решава кои нишки изпълняват тези future-и. Въпреки това се появява същата идея за собственост: една задача трябва да притежава това, което пази по време на изпълнението си, или да получава заемания, които остават валидни, докато приключи. Затова е важно първо да разбереш move, макар крайният код да използва Tokio.
Arc, Mutex и данната, която живее в ключалката
Един Rc<T> позволява на няколко собственици в една-единствена нишка да споделят стойност. Броячът му на референции не е атомарен, така че не може да се споделя между нишки. Arc<T> означава atomic reference counted: върши същата обща функция, но обновява брояча на референциите безопасно между нишки. Клонирането на Arc не клонира T; само създава още един собственик на същото заделяне.
Да имаш няколко собственици не е същото като да имаш разрешение да променяш. Ако T е променливо и няколко задачи могат да го достъпват, е нужна координация. Mutex<T> съдържа данната и позволява само на една задача в даден момент да получи променлив достъп чрез lock(). Резултатът от lock() е MutexGuard<T>. Докато този guard съществува, ключалката остава заета; когато излезе от областта си, неговият Drop освобождава ключалката. Не ти трябва отделно извикване на unlock, и това намалява риска да забравиш да освободиш ресурса при някой път на връщане.
Фиг. 7.2 | Данна, споделена между нишки.
// fig07_02.rs
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::thread;
fn main() {
let estado = "OK".to_string();
let estados = Arc::new(Mutex::new(HashMap::new()));
let copia = Arc::clone(&estados);
let h = thread::spawn(move || {
copia.lock().unwrap().insert("x".to_string(), estado);
});
h.join().unwrap();
println!("{:?}", estados.lock().unwrap());
}
$ rustc --edition 2024 fig07_02.rs && ./fig07_02
{"x": "OK"}
Формата Arc<Mutex<HashMap<_, _>>> се чете отвътре навън. HashMap е данната. Mutex е единствената врата, за да се промени. Arc позволява на различни нишки да са собственици на тази врата. За да вмъкне, нишката взема ключалката и получава guard-а; за да отпечата, main взема отново ключалката. Никога не се появява променлива референция към картата, без да съществува guard-ът.
В Go е честа практика да се декларира sync.Mutex до картата и да се следва конвенцията да се извиква Lock, преди да се пипне. Тази конвенция може да се капсулира добре, но езикът не задължава картата да е физически вътре в mutex-а. В Rust, като обвиеш данната в Mutex<T>, нормалният API не позволява да се получи &mut T без MutexGuard. Това не прави невъзможни всички грешки на конкурентността: все още можеш да предизвикаш мъртва хватка (deadlock), като вземаш ключалки в несъвместим ред, или да държиш guard твърде дълго. Прави обаче невъзможно в безопасен код състезание за данни, причинено от едновременно заемане на една и съща променлива стойност без синхронизация.
Не използвай lock().unwrap() като формула, за която не се мисли. lock() може да върне грешка, ако друга нишка е направила panic!, докато е държала ключалката; това се нарича отравяне (poisoning) на mutex-а. В един педагогически пример unwrap() прави този случай видим с panic!. В една реална услуга трябва да решиш дали това състояние обезсилва програмата, дали можеш да възстановиш данната с into_inner, или дали е удачно да преработиш дизайна, така че споделеното състояние да не е нужно.
revisor избягва Mutex<HashMap<...>>, защото не трябва да пълни споделена карта, докато пристига всеки отговор. Всеки future произвежда собствения си Estado, а join_all ги събира. Ресурсът, който наистина споделя, е лимит на редове: няколко future-а трябва да поискат разрешение да започнат запитване, не да променят обща колекция. Затова проектът използва Arc<Semaphore>, а не Arc<Mutex<Vec<Estado>>>.
Канали: предаване на стойности вместо споделяне на колекция
Каналът разделя комуникацията на два края: изпращач, Sender<T>, и получател, Receiver<T>. Изпращачът предава стойности с send; получателят ги взема с recv или като итерира върху него. Вместо няколко нишки да пишат в споделена структура, всеки производител предава стойност, която става собственост на получателя. Тази архитектура намалява споделената зона и изяснява кой сглобява крайния резултат.
Каналът, създаден с std::sync::mpsc::channel(), е с множество производители и един потребител. „Множество“ е, защото можеш да клонираш изпращача, преди да преместиш копие във всяка нишка. Получателят не се клонира: едно-единствено място решава какво да прави с всяко съобщение. Когато всички изпращачи изчезнат, итерирането върху получателя завършва. Това затваряне на канала е част от протокола, не случайна подробност.
Фиг. 7.3 | Канал предава състояния на нишката, която отпечатва доклада.
// fig07_03.rs
use std::sync::mpsc;
use std::thread;
fn main() {
let (emisor, receptor) = mpsc::channel();
let hilo = thread::spawn(move || {
for estado in ["catalogo: OK", "pagos: OK"] {
emisor.send(estado).unwrap();
}
});
hilo.join().unwrap();
for estado in receptor {
println!("{estado}");
}
}
$ rustc --edition 2024 fig07_03.rs && ./fig07_03
catalogo: OK
pagos: OK
Програмата чака нишката, преди да обходи получателя, за да запази детерминиран изход. В програма, която наистина обработва резултати, докато пристигат, обикновено би започнал да получаваш, докато производителите още работят. Тогава редът би бил този на пристигане, не непременно на входните услуги. Това може да е добро решение за интерфейс, който съобщава напредък, но не и за доклад, който трябва да подравни всяко състояние с оригиналната услуга.
Каналите не заместват автоматично mutex-ите. Ако няколко задачи трябва да четат и обновяват една и съща сметка, може би Mutex е естественият модел. Ако една задача произвежда стойности, а друга решава как да ги съхрани или покаже, каналът обикновено представя по-добре отговорността. Честата грешка е да се избира по мода: „mutex-ите са лоши“ или „каналите са сложни“. Полезният въпрос е кой трябва да е собственик на всяка данна във всеки момент.
В revisor join_all изпълнява сходна функция със получаването на всички резултати, но с допълнителен договор: запазва реда на входните future-и, така че estados[i] винаги е резултатът от servicios[i]. Това е, от което докладът има нужда: да знае на коя услуга принадлежи всяко състояние. Внимание с това, което не обещава: докладът не се отпечатва в реда на YAML списъка, защото reporte::tabla и reporte::json подреждат редовете по име на услуга, преди да ги запишат (ще го видиш в урок 8). Това, което join_all гарантира, е съответствието между всяка услуга и нейното състояние, а с канал, където резултатите пристигат в реда, в който приключват, би трябвало да се възстановява на ръка. Ако продуктът трябваше да отпечатва „приключи pagos“ веднага щом пристигне отговорът, канал или stream (последователност от стойности, които пристигат с времето) би била разумна опция. Не го добавяй само защото съществува; сегашният избор е умишлен и запазва доклада възпроизводим.
Send, Sync и грешката, която компилаторът спира
Send и Sync са маркерни trait-ове. Обикновено не изискват собствени методи; описват свойства за безопасност, които Rust може да изведе от полетата на един тип. Тип Send може да се прехвърли по стойност в друга нишка. Тип Sync може да се споделя по референция между нишки: ако T е Sync, то &T е Send. Много обичайни структури ги имплементират автоматично, когато компонентите им също са безопасни, но Rc<T> не е нито Send, нито Sync, защото броячът му не може да се обновява от няколко нишки.
Това правило не е списък за запаметяване. Това е въпрос, на който Rust отговаря чрез композиция. Ако направиш struct, който съдържа Rc<RefCell<_>>, той наследява ограниченията на тези части. Ако преминеш към Arc<Mutex<_>>, променяш представянето и също така наличните гаранции. Компилаторът следва стойността до closure-а, който се изпраща на thread::spawn, и изисква границата да е безопасна.
В Go едно състезание за данни може да се компилира и да изисква go test -race, за да се открие при изпълнение, което стигне точно до проблемното преплитане. Детекторът е ценен и трябва да го използваш, но зависи от това тестът да изпълни конфликтния път. Rust избягва състезанията за данни в безопасен код преди да се изпълни програмата. Това не доказва, че логиката е правилна, нито открива автоматично мъртви хватки, гладуване (starvation: една задача, която никога не получава ред, защото други заграбват ресурса) или зле проектирани протоколи. Съществува и unsafe, където програмистът поема допълнителни отговорности. Точното твърдение е: Rust избягва състезанията за данни чрез правилата си за типове и заемания в безопасен код; не обещава, че всяка конкурентна програма е правилна.
Вътре в revisor Arc<Semaphore> е валиден, защото семафорът на Tokio е проектиран да се споделя между задачи. Всеки future получава собствен Arc, иска разрешение и пази това разрешение по време на заявката. Типът на разрешението и жизненият му цикъл изразяват, че редът не може да се върне, преди да приключи запитването. Няма споделен брояч usize, който всеки future да увеличава и намалява ръчно.
async, future-и и runtime-ът на Tokio
Функция, означена с async fn, не изпълнява веднага цялото си тяло при извикване. Произвежда future: стойност, която представя изчакваща работа. Този future напредва, когато един изпълнител (executor) го poll-ва. Ако стигне до операция, която още не е готова, като чакане на мрежов отговор, връща контрола на изпълнителя. По-късно, когато операцията може да продължи, изпълнителят го poll-ва отново.
.await е точката, където една async функция чака резултата на друг future. Само по себе си не създава нова задача нито нова нишка. Това разграничение поправя две честни недоразумения. Първо: писането на let futuro = revisar(...); не стартира непременно заявката; само строи future-а. Второ: извикването на .await една след друга в същия блок може да направи операциите последователни. За да стартираш няколко операции конкурентно, строиш няколко future-а и ги водиш заедно с комбинатор като join_all, или ги превръщаш в задачи с tokio::spawn, когато наистина имаш нужда от независимост.
Rust дефинира синтаксиса и trait-овете на async, но не включва пълен async изпълнител в стандартната библиотека. Tokio е runtime-ът, избран от този проект. Предоставя изпълнител, таймери, async синхронизация и адаптации за вход-изход. Някои async библиотеки са независими от runtime-а, но ресурсите на Tokio, като таймерите му и няколко типа синхронизация, изискват да се изпълняват в контекст на Tokio. Прочети документацията на всеки crate, преди да предположиш, че всеки future работи еднакво с всеки runtime.
#[tokio::main]
async fn main() -> ExitCode {
let args = Args::parse();
ejecutar(&args).await
}
async fn ejecutar(args: &Args) -> ExitCode {
Атрибутът #[tokio::main] строи runtime-а и изпълнява async функцията main. main на програмата продължава да връща ExitCode, както видя в урок 4; променя се това, че вече може да чака async операции, преди да реши изходния код. Бинарният файл запазва отговорността да парсва аргументи, да отпечатва и да излиза; библиотеката запазва логиката за питане на услугите.
Не блокирай нишка на runtime-а със std::thread::sleep, тежко четене на файлове или дълго изчисление вътре в async функция. Блокирана нишка не може да poll-ва други future-и, възложени на нея. За блокираща работа съществува tokio::task::spawn_blocking; за мрежов вход-изход използвай async API-та като reqwest. revisor използва reqwest::Client и чака send().await, така че докато един отговор е в очакване, runtime-ът може да придвижва запитвания към други услуги.
Преди да видиш как revisor използва Tokio, е добре да видиш Tokio сам. Следващата програма е най-малката, която показва важното: три „проверки“, които вместо да питат мрежата, просто чакат, всяка различно време. Не се побира в чист rustc, защото зависи от crate-овете tokio и futures; затова живее в programas/revisor/examples/ и се изпълнява с Cargo, който изтегля и компилира тези зависимости.
Пример за cargo с tokio | Три чакания, водени едновременно с join_all: приключват в един ред и се предават в друг.
// ejemplo_tokio.rs
use std::time::{Duration, Instant};
use futures::future::join_all;
use tokio::time::sleep;
async fn revisar(nombre: &str, espera_ms: u64) -> String {
sleep(Duration::from_millis(espera_ms)).await;
println!("terminó {nombre}");
format!("{nombre}: respondió tras {espera_ms} ms")
}
#[tokio::main]
async fn main() {
let servicios = [("catalogo", 600), ("pagos", 200), ("usuarios", 400)];
let inicio = Instant::now();
let futuros = servicios.iter().map(|(nombre, ms)| revisar(nombre, *ms));
let resultados = join_all(futuros).await;
println!("--- en el orden de la lista ---");
for resultado in &resultados {
println!("{resultado}");
}
// Esperarlos uno tras otro habría tardado 1200 ms; a la vez tardan lo del más lento.
let a_la_vez = inicio.elapsed() < Duration::from_millis(1100);
println!("tardó menos que la suma de las esperas: {a_la_vez}");
}
$ cargo run --example ejemplo_tokio
terminó pagos
terminó usuarios
terminó catalogo
--- en el orden de la lista ---
catalogo: respondió tras 600 ms
pagos: respondió tras 200 ms
usuarios: respondió tras 400 ms
tardó menos que la suma de las esperas: true
Изпълни го от programas/revisor/ (при първото изпълнение Cargo се бави известно време да компилира зависимостите).
#[tokio::main] превръща main в async функция: строи runtime-а и му предава future-а, който main описва. tokio::time::sleep е чакането на Tokio и прилича на std::thread::sleep по това, което прави, но не по начина, по който го прави: с .await задачата отстъпва контрола на runtime-а, докато чака, а runtime-ът се възползва, за да придвижи останалите. Със std::thread::sleep цялата нишка би заспала и нищо друго не би напредвало в нея.
Чети изхода на две части. Съобщенията terminó ... излизат в реда, в който се изпълнява всяко чакане: pagos (200 ms), usuarios (400 ms) и catalogo (600 ms), въпреки че списъкът ги декларира в друг ред. След това join_all предава резултатите в реда на списъка —catalogo, pagos, usuarios—, защото връща Vec, където всяка позиция съответства на входния си future. Последният ред проверява, че е било конкурентно: чакането на трите проверки една след друга би отнело 1200 ms, а вместо това отнема колкото най-бавната, около 600 ms.
Ако смениш sleep с мрежово извикване с .send().await, имаш формата на revisor: много future-и, които чакат, един-единствен join_all, който ги води, и Vec от резултати, подравнен със списъка на услугите.
join_all, семафори и лимитът на паралелизма в revisor
Да пуснеш всички възможни заявки едновременно не винаги е подобрение. Файл с хиляди услуги би могъл да отвори твърде много връзки, да претовари локалната мрежа, да изчерпи файловите дескриптори или да натовари сървъра, който точно се опитваш да проверяваш. Конкурентността се нуждае от лимит. Аргументът --paralelo на revisor изразява колко запитвания могат да са активни едновременно.
Един семафор съдържа разрешения. За да започне запитване, един future придобива едно; ако не са останали, чака. Когато разрешението излезе от областта си, се освобождава автоматично и друг future може да продължи. Това е същата идея на RAII (ресурсът се освобождава, когато стойността, която го представя, излезе от областта си), която видя с MutexGuard: ресурсът се освобождава при унищожаването на guard-а, дори ако функцията излезе по нормален път. Тук ресурсът не е изключителна ключалка, а ограничен капацитет.
pub async fn revisar_todos(
cliente: &reqwest::Client,
servicios: &[Servicio],
paralelo: usize,
) -> Vec<Estado> {
let paralelo = if paralelo == 0 {
PARALELO_POR_OMISION
} else {
paralelo
};
let turnos = Arc::new(Semaphore::new(paralelo));
let futuros = servicios.iter().map(|s| {
let turnos = Arc::clone(&turnos);
async move {
// pide turno; espera si ya hay `paralelo` corriendo, y lo devuelve al soltar `_turno`
let _turno = turnos.acquire().await.expect("el semáforo nunca se cierra");
revisar(cliente, s).await
}
});
join_all(futuros).await
}
servicios.iter() запазва заеманията към списъка; не консумира услугите. Всеки async closure притежава своя клон на Arc<Semaphore>, но взема назаем cliente и s по време на извикването на revisar. Това работи, защото join_all(futuros).await приключва, преди revisar_todos да може да се върне, така че тези заемания остават живи. Ако използваше tokio::spawn, задачата би могла да надживее функцията, която я е създала, и обикновено би се нуждаела от данни с живот 'static; там би трябвало да преместиш или клонираш повече данни.
join_all връща Vec<Estado> в реда на входните future-и, не в реда, в който заявките приключват. Това е полезно решение за доклада: състояние нула съответства на услуга нула. Семафорът ограничава кога всеки future може да влезе в revisar, но не променя тази крайна връзка. Така програмата има истинска конкурентност, без да превръща доклада в източник на случаен ред.
Функцията revisar не връща Result<Estado>. Това решение заслужава внимание. Един HTTP 500, един timeout или една отказана връзка не е вътрешна грешка, която пречи на програмата да продължи: е точно информацията, която revisor трябва да докладва за една услуга. Затова тези случаи се превръщат във варианти Estado::Falla. Грешката на acquire се третира различно, защото семафорът никога не се затваря в този дизайн; ако се случи, би била нарушение на вътрешно предположение.
Timeout-ът на проекта се конфигурира върху заявката на reqwest, с timeout_ms на всяко Servicio. Не бъркай този лимит със семафора. Timeout-ът ограничава колко може да чака едно отделно запитване; семафорът ограничава колко запитвания могат да чакат или да комуникират едновременно. Имаш нужда от двете: без timeout един ред може да остане зает твърде дълго; без семафор много заявки с timeout могат да стартират всички заедно и да претоварят ресурси.
Честно сравнение с Go
Go прави много лесно стартирането на конкурентна единица: go revisar(s). Runtime-ът планира goroutine-и върху нишки на операционната система, разширява стековете им и управлява мрежовата работа. Rust разделя решението изрично: thread::spawn създава нишка на операционната система; runtime като Tokio изпълнява async future-и; tokio::spawn създава задача на Tokio. Това означава, че Go обикновено има по-малко церемония в началото, а Rust те принуждава да знаеш коя от трите абстракции използваш.
Rust няма магическо предимство в производителността от писането на async. Една async задача не ускорява отделно HTTP запитване; подобрява използването на нишките, докато няколко запитвания чакат. При натоварване на процесора async може да е по-лош, ако блокира изпълнителя. В този случай използвай нишки, пул от worker-и или spawn_blocking. Go също не прави автоматично по-бързо едно интензивно за процесора изчисление: няколко goroutine-и могат да се състезават за едни и същи ядра. Да измериш типа работа има по-голямо значение, отколкото да приложиш модна дума.
Централната гаранция на Rust се появява, преди да се изпълни: една променлива данна не може да се заема по несъвместим начин, а стойностите, изпратени към друга нишка, трябва да отговарят на подходящите trait-ове. Go предпочита малък синтаксис и инструменти по време на изпълнение, като детектора на състезания. Go може правилно да капсулира mutex-и и канали; Rust може да има deadlock-ове и логически грешки. Разликата не е „Go допуска грешки, а Rust не“. Тя е къде всеки език поставя тежестта на проверката и какви грешки може да отхвърли, преди да се изпълни.
За revisor решението е оправдано от проблема. Има много HTTP чакания, всеки резултат трябва да остане свързан със своята услуга и е нужен конфигурируем таван на заявките. Tokio, join_all и Semaphore изразяват тези три нужди. Дизайн с една нишка на услуга би работил за малък списък, но би мащабирал по-зле и не носи предимство за случая на употреба. Дизайн с mutex и споделена карта също би могъл да работи, но би усложнил връзка, която join_all вече запазва.
Грешката, която ще видиш
Най-поучителната грешка на този урок е E0277. Не означава, че „Rust не иска да използва нишки“; означава, че типът, който се опитваш да преместиш, не удовлетворява договора, който изисква thread::spawn. Rc<i32> е полезен за споделяне на собственост в рамките на една нишка, но броячът му на референции не е атомарен. Преместването му в closure, който може да се изпълни в друга нишка, би било небезопасно.
Фиг. 7.4 | Rc<T> не може да се изпрати към нишка.
// fig07_04.rs
use std::rc::Rc;
use std::thread;
fn main() {
let conteo = Rc::new(0);
thread::spawn(move || println!("{conteo}"));
}
$ rustc --edition 2024 fig07_04.rs
error[E0277]: `Rc<i32>` cannot be sent between threads safely
--> fig07_04.rs:7:19
|
7 | thread::spawn(move || println!("{conteo}"));
| ------------- -------^^^^^^^^^^^^^^^^^^^^^
| | |
| | `Rc<i32>` cannot be sent between threads safely
| | within this `{closure@fig07_04.rs:7:19: 7:26}`
| required by a bound introduced by this call
|
= help: within `{closure@fig07_04.rs:7:19: 7:26}`, the trait `Send` is not implemented for `Rc<i32>`
note: required because it's used within this closure
--> fig07_04.rs:7:19
|
7 | thread::spawn(move || println!("{conteo}"));
| ^^^^^^^
note: required by a bound in `spawn`
--> /rustc/48a229ceaefd4985c50990b14116b6d856af0985/library/std/src/thread/functions.rs:125:0
error: aborting due to 1 previous error
For more information about this error, try `rustc --explain E0277`.
Съобщението съдържа отговора. Closure-ът улавя Rc<i32>, thread::spawn изисква уловеното да е Send, а Rc<i32> не имплементира Send. Не поправяй тази грешка, като добавяш trait-ове на ръка с unsafe impl Send; би обещал на компилатора безопасност, която Rc не предлага. Ако няколко нишки трябва само да четат една данна, използвай Arc<T>. Ако трябва и да я променят, прецени Arc<Mutex<T>> или преработи потока, за да изпращаш стойности по канали.
Разпознай и грешката в дизайна, която понякога компилаторът не отхвърля: да държиш MutexGuard по време на операция .await. Ако стартираш задачата с tokio::spawn и guard-ът е от std::sync::Mutex, rustc наистина я отхвърля (future cannot be sent between threads safely, защото този guard не е Send); но ако future-ът се чака в същата нишка, както прави join_all в revisor, програмата се компилира. За този случай clippy носи по подразбиране предупреждението await_holding_lock. Guard на std::sync::Mutex блокира нишка; guard на async mutex пази ключалката, докато задачата може да отстъпи изпълнителя. И в двата случая чакането на мрежа, докато държиш ключалката, обикновено блокира работа без нужда и може да доведе до мъртви хватки. Извлечи или обнови данната под ключалката, освободи guard-а и чак тогава чакай.
Какво се прави погрешно
Създаване на нишка за всяка услуга без лимит. Работи с три примера и се проваля като стратегия, когато списъкът расте. Нишките на операционната система консумират ресурси на операционната система и голямо количество от тях прави по-трудно планирането, измерването и отстраняването на грешки. За мрежов вход-изход използвай async с лимит на паралелизма; за процесор използвай брой нишки, пропорционален на работата и на наличните ядра.
Използване на
Arc<Mutex<_>>като автоматичен отговор на всяка грешка на ownership. Тази комбинация е правилна, когато има наистина споделено променливо състояние, но може да скрие дизайн, в който няколко задачи правят твърде много. Ако всяка задача може да върне стойност и само една част ги събира, канал,join_allили последваща редукция обикновено е по-ясно и намалява конкуренцията за ресурс (contention).Държане на ключалка по време на мрежова заявка или
.await. Guard-ът съществува, за да защитава малка критична секция. Ако го запазиш, докато чакаш, превръщаш конкурентните задачи в опашка и увеличаваш възможността за мъртва хватка. Ограничи областта му с фигурни скоби или с временна променлива, така че guard-ът да се унищожи преди чакането.Извикване на блокиращи функции в код на Tokio.
std::thread::sleepблокира нишката, не само една задача. Голямо четене, блокираща заявка или тежко изчисление имат същия проблем. Използвай async API-та за вход-изход,tokio::timeза таймери илиspawn_blockingза работа, която наистина трябва да блокира.Бъркане на
asyncс паралелизъм. Един future може да напредва конкурентно с други и все пак да се изпълнява върху една-единствена нишка. Ако трябва да ускориш изчисление на процесора, трябва да решиш как да го разпределиш между ядрата. Ако чакаш мрежа, async подобрява използването на съществуващите нишки. Преди да оптимизираш, измери какво чака програмата.Използване на
tokio::spawnсамо за да „стане конкурентно“. Вrevisorjoin_allможе да води future-и, които заематclienteиservicios, запазва реда и избягва да изисква собственост'static.tokio::spawnе полезен за независими задачи, които трябва да живеят отвъд текущия блок, но предполага друг договор за живот и за типове.Третиране на изхода, който пристига първи, като състоянието на услугата, която заема тази позиция. В интерфейс за напредък може да е полезно да се съобщава по пристигане. В доклад, ако резултатът на друга услуга се промъкне в един ред, всеки ред лъже за своята услуга.
revisorизбягва този риск сjoin_all, който запазва съответствието междуservicios[i]иestados[i], а интеграционните му тестове проверяват това свойство; редът, в който се отпечатват редовете, е друго решение, което взема докладът, като ги подрежда по име.
Упражнения
Упражнение 1 — Три проверки с join
Напиши програма с три текстови Servicio. Използвай thread::spawn(move || ...), за да произведеш по едно състояние за всяка, запази handle-овете във вектор и използвай join, за да отпечаташ резултатите в същия ред като входа. Не използвай sleep, за да „дадеш време“ на нишките.
Упражнение 2 — Защитен брояч
Създай Arc<Mutex<u32>> с начална стойност нула. Стартирай четири нишки; всяка трябва да увеличи брояча веднъж. Изчакай всички handle-ове и отпечатай total: 4. После смени нарочно Arc с Rc и потвърди, че се появява E0277.
Упражнение 3 — Резултати по канал
Създай канал и три нишки производители. Всеки производител трябва да изпрати името на една услуга и едно състояние. Главната нишка трябва да получи точно три съобщения и да ги подреди, преди да ги отпечата, по име. Обясни в един коментар защо не бива да зависиш от реда на пристигане.
Упражнение 4 — Обясни лимита на revisor
Прочети programas/revisor/src/revisar.rs. Изпълни интеграционните тестове и намери теста el_tope_de_paralelo_se_respeta. В дневника си отговори: какъв ресурс контролира Semaphore, кога се придобива разрешението, кога се освобождава и защо join_all запазва реда на servicios.
Решения
Решение 1
Правилното решение премества всяко Servicio в нишката и събира handle-овете след това. Решаващото не е map, а това, че векторът от handle-ове поддържа задължението да се изчака всяка работа, преди да приключи main. Ако отпечатваш след всеки join, наблюдаваният ред е редът на вектора от handle-ове.
Един начин да провериш, че не си зависил от sleep, е да изпълниш програмата няколко пъти: трябва винаги да отпечатва трите реда. Планировчикът може да промени коя нишка приключва първа, но не може да попречи на join да чака.
Решение 2
Всяка нишка трябва да получи собствен Arc::clone(&contador). Вътре в нишката вземи guard-а, увеличи и остави guard-а да излезе от областта си. След като изчакаш четирите handle-а, вземи един последен guard, за да отпечаташ стойността. Не се опитвай да държиш променливо заемане на брояча извън mutex-а; това заемане не може да съжителства с останалите нишки.
При смяната на Arc с Rc очакваният резултат не е числов изход, а E0277. Поправката не се състои в „заглушаване“ на компилатора: Rc служи за споделени референции в една нишка; Arc е подходящият тип, за да има броячът собственици в няколко нишки.
Решение 3
Всеки производител получава клон на изпращача и изпраща структура или кортеж с името и състоянието. Оригиналният изпращач трябва да престане да съществува, преди да обходиш целия получател, или можеш да извикаш recv точно три пъти, защото знаеш броя на производителите. Запази получените съобщения във вектор и го подреди по име, преди да отпечаташ.
Важната част е, че получателят е единственият собственик на крайния вектор. Производителите нямат променлив достъп до него. Затова не ти трябва mutex, за да събереш резултатите; собствеността пътува по канала заедно с всяко съобщение.
Решение 4
Семафорът контролира броя на HTTP запитванията, които могат да са активни едновременно, не общия брой услуги. Всеки future придобива разрешение точно преди да извика revisar(cliente, s).await. Разрешението живее в _turno; когато това извикване приключи, _turno излиза от областта си и връща капацитета на семафора.
join_all получава future-и, построени при обхождането на servicios, и произвежда вектора от състояния в същата последователност. Затова резултатите могат да приключват в различно време, без да се разминават с услугата, която ги е породила. Тестът измерва времена, за да потвърди, че лимитът променя поведението и не е само декоративен флаг.
Как да разбера, че съм успял
Изпълни фигурите на този урок от техните директории и сравни точния изход:
cd programas/07-concurrencia-async
rustc --edition 2024 fig07_01.rs && ./fig07_01
rustc --edition 2024 fig07_02.rs && ./fig07_02
rustc --edition 2024 fig07_03.rs && ./fig07_03
rustc --edition 2024 fig07_04.rs
Първите три компилации трябва да завършат без предупреждения и да произведат документираните изходи. Последната трябва да се провали с E0277; този провал е правилният резултат от упражнението.
Провери, че блоковете и техните изходи остават проверими от корена на курса:
herramientas/verificar-programas.sh es
herramientas/verificar-extractos.sh
herramientas/verificar-ejemplos.sh
Трите команди трябва да завършат успешно. Първата потвърждава, че всяка фигура се компилира, изпълнява и съвпада с документирания си изход, освен фигурата, проектирана да се провали. Втората потвърждава, че откъсите от revisor остават точни копия на реалния проект. Третата компилира и изпълнява примера с tokio от този урок и сравнява това, което отпечатва, с документираното.
Накрая провери реалното конкурентно поведение на проекта:
cd programas/revisor
cargo test
cargo clippy --all-targets -- -D warnings
cargo fmt --check
cargo test трябва да съобщи успешни резултати, включително теста el_tope_de_paralelo_se_respeta. cargo clippy --all-targets -- -D warnings трябва да завърши без предупреждения, а cargo fmt --check не бива да предлага промени. Ако можеш да обясниш защо семафорът ограничава запитванията, защо join_all запазва реда и защо Rc предизвиква E0277, завърши урока.
За по-нататъшно четене
The Rust Programming Language, глава 16: Fearless Concurrency — консултирано на 2 октомври 2026 г.
The Rust Programming Language, глава 17: Fundamentals of Asynchronous Programming — консултирано на 2 октомври 2026 г.
Документация на
std::thread— консултирано на 2 октомври 2026 г.Документация на
tokio::sync::Semaphore— консултирано на 2 октомври 2026 г.
Предпочитате имейл? Пишете ни на hola@habil.mx