Canais (channels)
Nesta aula, exploramos canais (channels) em Rust, focando no módulo mpsc (múltiplos produtores, único consumidor). Aprendemos a usar send e recv para comunicação entre threads, lidar com múltiplos produtores e aplicar padrões como roteamento de mensagens e sincronização.
Nesta aula, vamos mergulhar nos canais (channels) em Rust, uma ferramenta essencial para comunicação segura entre threads. Canais seguem o modelo de passagem de mensagens, onde um produtor envia dados e um consumidor os recebe, evitando problemas de concorrência como data races. O módulo std::sync::mpsc implementa canais com múltiplos produtores e único consumidor (multiple producer, single consumer), sendo a base para a maioria dos casos de uso.
Entender canais é fundamental para escrever programas concorrentes em Rust. Eles permitem que threads troquem informações de forma controlada, sem compartilhamento de estado, alinhando-se com a filosofia de segurança da linguagem. Veremos desde o básico até padrões mais avançados.
mpsc
O módulo std::sync::mpsc fornece canais assíncronos (com buffer) e síncronos (sem buffer, via SyncSender). A função channel() retorna um par (Sender, Receiver). O Sender pode ser clonado para múltiplos produtores, enquanto Receiver é único. O canal usa um buffer interno (atualmente com capacidade ilimitada) para armazenar mensagens até que sejam lidas.
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
tx.send(String::from("Olá")).unwrap();
});
let received = rx.recv().unwrap();
println!("Recebido: {}", received);
}O método send retorna Result<(), SendError>; se o receptor foi descartado, retorna erro. recv bloqueia até receber uma mensagem ou o canal fechar (todos os produtores descartados), retornando Err se não houver mais produtores.
send e recv
send envia um valor pelo canal. Se o receptor já foi dropado, retorna SendError contendo o valor de volta. recv bloqueia a thread atual até que uma mensagem esteja disponível ou o canal seja fechado. Há também try_recv, que não bloqueia e retorna TryRecvError se vazio ou fechado.
use std::sync::mpsc;
use std::thread;
use std::time::Duration;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
thread::sleep(Duration::from_secs(1));
tx.send("mensagem").unwrap();
});
match rx.recv_timeout(Duration::from_secs(2)) {
Ok(msg) => println!("Recebido: {}", msg),
Err(e) => println!("Erro: {:?}", e),
}
}O método recv_timeout bloqueia por um tempo máximo. Também é possível iterar sobre o receptor com um loop for, que termina quando o canal é fechado.
Múltiplos produtores
Para ter múltiplos produtores, clonamos o Sender. Cada clone pode ser movido para uma thread diferente. O Receiver permanece único e coleta todas as mensagens na ordem em que foram enviadas (não necessariamente FIFO estrito devido ao buffer, mas geralmente FIFO).
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
let tx1 = tx.clone();
thread::spawn(move || {
tx.send("do produtor 1").unwrap();
});
thread::spawn(move || {
tx1.send("do produtor 2").unwrap();
});
for received in rx {
println!("Recebido: {}", received);
}
}Note que o loop for sobre o receptor termina quando todos os Senders (incluindo clones) forem dropados. Isso é útil para processar um fluxo finito de mensagens.
Padrões
Padrões comuns com canais incluem: roteamento de mensagens (usando enum para diferentes tipos), sincronização (enviar sinal de conclusão), e balanceamento de carga (vários workers consumindo de um canal).
enum Mensagem {
Trabalho(String),
Sair,
}
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel::<Mensagem>();
let worker = thread::spawn(move || {
loop {
match rx.recv().unwrap() {
Mensagem::Trabalho(t) => println!("Processando: {}", t),
Mensagem::Sair => break,
}
}
});
tx.send(Mensagem::Trabalho("tarefa1".into())).unwrap();
tx.send(Mensagem::Trabalho("tarefa2".into())).unwrap();
tx.send(Mensagem::Sair).unwrap();
worker.join().unwrap();
}Outro padrão é usar canais para sincronizar o início de threads (barreira). Enviar um valor vazio como sinal.
Boas práticas
- Sempre trate erros de
senderecv; o canal pode fechar inesperadamente. - Use
try_recvem loops não bloqueantes. - Clone o
Senderapenas quando necessário; cada clone aumenta a contagem de produtores. - Para comunicação bidirecional, crie dois canais.
Referências
- Documentação oficial std::sync::mpsc
- The Rust Book: Message Passing
- Rust Design Patterns: Concurrency
- Rust by Example: Channels
- The Rustonomicon: Arc and Channels
Exercícios
Crie um programa que usa um canal para enviar 10 números inteiros de uma thread para a main. A thread produtora deve enviar os números de 1 a 10, e a main deve imprimi-los.
✓ Resposta:use std::sync::mpsc; use std::thread; fn main() { let (tx, rx) = mpsc::channel(); thread::spawn(move || { for i in 1..=10 { tx.send(i).unwrap(); } }); for received in rx { println!("{} ", received); } }Modifique o exercício anterior para usar dois produtores: um envia números pares e outro ímpares. O consumidor (main) deve imprimir todos.
✓ Resposta:use std::sync::mpsc; use std::thread; fn main() { let (tx, rx) = mpsc::channel(); let tx1 = tx.clone(); thread::spawn(move || { for i in (2..=10).step_by(2) { tx.send(i).unwrap(); } }); thread::spawn(move || { for i in (1..=9).step_by(2) { tx1.send(i).unwrap(); } }); for received in rx { println!("{} ", received); } }Escreva um programa que usa um enum de mensagens para enviar comandos:
Texto(String),Numero(i32)eSair. Uma worker thread processa as mensagens até receberSair.✓ Resposta:use std::sync::mpsc; use std::thread; enum Mensagem { Texto(String), Numero(i32), Sair, } fn main() { let (tx, rx) = mpsc::channel::<Mensagem>(); let worker = thread::spawn(move || { loop { match rx.recv().unwrap() { Mensagem::Texto(t) => println!("Texto: {}", t), Mensagem::Numero(n) => println!("Número: {}", n), Mensagem::Sair => break, } } }); tx.send(Mensagem::Texto("Olá".into())).unwrap(); tx.send(Mensagem::Numero(42)).unwrap(); tx.send(Mensagem::Sair).unwrap(); worker.join().unwrap(); }Use
recv_timeoutpara ler mensagens com timeout de 1 segundo. Se não houver mensagem, imprima "timeout".✓ Resposta:use std::sync::mpsc; use std::thread; use std::time::Duration; fn main() { let (tx, rx) = mpsc::channel(); thread::spawn(move || { thread::sleep(Duration::from_secs(2)); tx.send("depois de 2s").unwrap(); }); loop { match rx.recv_timeout(Duration::from_secs(1)) { Ok(msg) => { println!("Recebido: {}", msg); break; } Err(mpsc::RecvTimeoutError::Timeout) => println!("timeout"), Err(mpsc::RecvTimeoutError::Disconnected) => { println!("canal fechado"); break; } } } }Crie um programa que simula um pipeline: thread A gera números, envia para thread B que dobra e envia para thread C que imprime. Use dois canais.
✓ Resposta:use std::sync::mpsc; use std::thread; fn main() { let (tx_a, rx_a) = mpsc::channel(); let (tx_b, rx_b) = mpsc::channel(); // Thread A: gera números let produtor = thread::spawn(move || { for i in 1..=5 { tx_a.send(i).unwrap(); } }); // Thread B: dobra let processador = thread::spawn(move || { for num in rx_a { tx_b.send(num * 2).unwrap(); } }); // Thread C: imprime let consumidor = thread::spawn(move || { for num in rx_b { println!("Resultado: {}", num); } }); produtor.join().unwrap(); processador.join().unwrap(); consumidor.join().unwrap(); }