Stream и каналы
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Stream и каналы
Последний урок блока. Футура это одно значение в будущем. А что, если значений много и они приходят по одному во времени? Это
Stream, асинхронный итератор. И связаны с ним каналы: как async-задачи передают друг другу данные, не блокируясь, и что делать, когда производитель быстрее потребителя. СоберёмStream, ограниченный канал с backpressure иoneshot, опираясь на свой рантайм.
Stream это асинхронный Iterator
Сравни две сигнатуры. У обычного итератора метод next(&mut self) -> Option<T>: вызвал, получил следующий элемент или None в конце. У Stream тот же смысл, но асинхронный:
pub trait Stream {
type Item;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>>;
}
Три исхода poll_next. Ready(Some(x)) это очередной элемент. Ready(None) это конец потока (как None у итератора). Pending это «элемента пока нет, разбужу через Waker». Это прямое скрещивание Iterator (Option на конце) и Future (Poll на готовности).
Поверх трейта удобно иметь метод next, который даёт футуру на один элемент, чтобы писать while let Some(x) = s.next().await:
pub trait StreamExt: Stream {
fn next(&mut self) -> Next<'_, Self> where Self: Unpin {
Next { stream: self }
}
}
impl<S: Stream + ?Sized> StreamExt for S {}
impl<S: Stream + Unpin + ?Sized> Future for Next<'_, S> {
type Output = Option<S::Item>;
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
Pin::new(&mut *self.stream).poll_next(cx)
}
}
Самый простой Stream оборачивает обычный итератор (он всегда готов, Pending не отдаёт):
let mut s = iter(vec![1, 2, 3]).map(|x| x * 10);
while let Some(x) = s.next().await {
println!("{x}"); // 10, 20, 30
}
Комбинаторы (map, filter, take) работают как у итераторов, только асинхронно: оборачивают исходный поток и преобразуют его poll_next. В реальном коде их даёт крейт futures (или tokio_stream), у нас они написаны руками в эталоне блока.
Каналы: как задачи передают данные
Каналы это основной способ общения async-задач. Tokio даёт четыре вида под разные задачи:
mpsc: много отправителей, один приёмник. Очередь сообщений, рабочая лошадка конвейеров.oneshot: ровно одно значение от одного к одному. Для возврата результата из задачи.broadcast: каждое сообщение видят все приёмники (fan-out).watch: одно всегда-актуальное значение, приёмники видят только последнее.
Соберём ограниченный mpsc сами, потому что в нём живёт важнейшее понятие: backpressure.
Backpressure
Backpressure это про то, что происходит, когда производитель быстрее потребителя. У неограниченного канала ответ плохой: очередь растёт, пока не съест всю память. У ограниченного канала (с ёмкостью) ответ правильный: когда буфер полон, send не проходит, а приостанавливается, пока потребитель не освободит место. Медленный потребитель так тормозит быстрого производителя, и система сама приходит в равновесие.
Вот ядро канала. send это футура: если место есть, кладёт значение и будит приёмника; если буфер полон, сохраняет Waker и возвращает Pending.
impl<T> Future for Send<'_, T> {
type Output = Result<(), SendError<T>>;
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
let this = unsafe { self.get_unchecked_mut() };
let mut inner = this.sender.inner.lock().unwrap();
if !inner.receiver_alive {
let value = this.value.take().unwrap();
return Poll::Ready(Err(SendError(value))); // приёмник ушёл
}
if inner.buffer.len() < inner.capacity {
let value = this.value.take().unwrap();
inner.buffer.push_back(value);
inner.wake_receiver(); // разбудили приёмника: появился элемент
Poll::Ready(Ok(()))
} else {
// Буфер полон: ждём, пока приёмник освободит место. Это backpressure.
inner.send_wakers.push(cx.waker().clone());
Poll::Pending
}
}
}
Приёмник симметричен и реализует Stream. Забирает элемент и будит отправителей, ждущих места; если буфер пуст и отправителей не осталось, поток закончился (Ready(None)):
impl<T> Stream for Receiver<T> {
type Item = T;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<T>> {
let mut inner = self.inner.lock().unwrap();
if let Some(value) = inner.buffer.pop_front() {
inner.wake_senders(); // освободили место: будим отправителей
Poll::Ready(Some(value))
} else if inner.senders == 0 {
Poll::Ready(None) // никого нет и пусто: конец потока
} else {
inner.recv_waker = Some(cx.waker().clone());
Poll::Pending
}
}
}
Заметь симметрию пробуждений: отправитель будит приёмника, когда положил элемент; приёмник будит отправителей, когда забрал и освободил место. Это та же механика «не звони мне, я позвоню тебе», только теперь между двумя задачами через общий буфер под Mutex. Ровно так устроен tokio::sync::mpsc.
Канал на ёмкости 1 наглядно показывает backpressure: второй send подряд не пройдёт, пока приёмник не заберёт первый.
let (tx, mut rx) = channel::<u32>(1);
spawner.spawn(async move {
for i in 0..3 {
tx.send(i).await.unwrap(); // второй send упрётся в backpressure
}
});
spawner.spawn(async move {
while let Some(x) = rx.next().await {
println!("получил {x}"); // 0, 1, 2 строго по одному
}
});
oneshot: одно значение и сигнал
oneshot это вырожденный случай: одно значение от одного к одному. Приёмник это сразу футура, без Stream. Отправитель берёт self по значению, поэтому послать можно ровно раз. Частый приём: вернуть результат из заспавненной задачи или передать разовый сигнал. А Drop отправителя без отправки будит приёмника с ошибкой, и это полезно: приёмник узнаёт, что ответа не будет, вместо вечного ожидания.
let (tx, rx) = oneshot::<&str>();
tx.send("привет").unwrap();
assert_eq!(rx.await, Ok("привет"));
let (tx, rx) = oneshot::<u32>();
drop(tx); // отправитель ушёл, не послав
assert_eq!(rx.await, Err(RecvError)); // приёмник не висит, а получает ошибку
Что унести из урока
Stream это асинхронный итератор: poll_next отдаёт Ready(Some(x)) на элемент, Ready(None) на конец, Pending на ожидание; комбинаторы map/filter работают как у обычных итераторов, только асинхронно. Каналы это способ общения задач: mpsc для конвейеров, oneshot для одного ответа, broadcast для fan-out, watch для актуального состояния. Главное понятие это backpressure: ограниченный канал тормозит быстрого производителя, заставляя send ждать места, и так бережёт память. Под капотом всё та же пара буфер плюс сохранённые Waker’ы обеих сторон, что мы собирали весь блок.
На этом блок «Async вглубь» закончен. Ты прошёл путь от async fn как машины состояний через Pin, свой исполнитель и reactor, внутренности Tokio до отмены и потоков. Async больше не магия: это poll, Waker и очередь задач, и ты собрал каждую деталь руками. Дальше в разделе ждёт многопоточность и атомики, где те же задачи параллельности решаются на уровне железа.