Источник, поток и приёмник
Основные строительные блоки потоковой обработки.
«Источник, поток и приёмник» — бесплатный урок Scala for Backend Engineering & Functional Programming на CoddyKit. Это урок 1 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения Scala for Backend Engineering & Functional Programming, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс Scala for Backend Engineering & Functional Programming содержит 4 уроков всего.
Три строительных блока
Akka Streams представляет конвейер обработки данных в виде графа этапов обработки. Три основных линейных этапа: Source (создаёт элементы), Flow (преобразует их) и Sink (потребляет их).
У источника один выход, у приёмника один вход, а у потока ровно один вход и один выход. Их соединение описывает, что должно произойти, но не когда.
Определение источника
Source[Out, Mat] выдаёт элементы типа Out и предоставляет материализованное значение типа Mat. Простейшие источники создаются из коллекций, находящихся в памяти, или диапазонов.
До запуска потока Source является лишь неизменяемым шаблоном, который можно свободно использовать повторно.
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")Определение приёмника
Sink[In, Mat] потребляет элементы типа In. Материализованное значение часто содержит результат потребления, например Future, завершающийся после окончания потока.
Sink.foreach выполняет побочный эффект для каждого элемента, а Sink.fold накапливает единый результат.
import akka.stream.scaladsl.Sink
import scala.concurrent.Future
val printSink: Sink[Int, Future[akka.Done]] =
Sink.foreach(println)
val sumSink: Sink[Int, Future[Int]] =
Sink.fold(0)(_ + _)Определение потока
Flow[In, Out, Mat] располагается между Source и Sink и преобразует каждый элемент. Потоки можно повторно использовать отдельно и компоновать до подключения к любой конечной точке.
В данном примере Flow удваивает целые числа и преобразует их в строки.
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")Соединение источника с приёмником
Оператор via подключает Flow к Source, а to подключает Sink. Прямое соединение Source с Sink через to создаёт RunnableGraph — замкнутый запускаемый шаблон.
Элементы пока не перемещаются; здесь лишь описывается топология.
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via: добавление потока
Используйте via, чтобы встроить Flow в конвейер. Вызов Source.via(flow) возвращает новый Source, тип выходных данных которого совпадает с выходным типом Flow.
Цепочка вызовов via позволяет строить длинные конвейеры преобразований из небольших, удобных для тестирования компонентов Flow.
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))Типобезопасность между этапами
Компилятор проверяет, что выходной тип каждого этапа совпадает с входным типом следующего. Нельзя подключить Source[Int] к Sink[String] без промежуточного Flow, преобразующего тип.
Эта статическая проверка обнаруживает ошибки соединения конвейера до его запуска.
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))Материализуемые значения
Каждая схема содержит материализуемое значение — дескриптор, создаваемый при запуске потока. Sink.fold материализует Future с результатом. По умолчанию при объединении этапов сохраняется крайнее левое материализуемое значение (NotUsed для обычных источников).
Используйте toMat и Keep, чтобы выбрать значение нужной стороны.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)Повторно используемые компоненты
Поскольку Sources, Flows и Sinks являются неизменяемыми значениями, их можно определить один раз и повторно использовать в разных конвейерах. Это способствует созданию библиотеки небольших именованных этапов обработки.
Flow, определённый для разбора данных, можно встроить и в конвейер обработки файла, и в конвейер обработки HTTP-запросов.
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)Распространённые конструкторы Source
Akka Streams предоставляет множество фабрик Source: Source.single, Source.repeat, Source.tick для выдачи данных по времени, Source.future из Future и Source.empty.
Правильный выбор конструктора явно показывает назначение источника данных.
import scala.concurrent.duration._
val ticks = Source.tick(0.seconds, 1.second, "tick")
val onceF = Source.future(scala.concurrent.Future.successful(42))
val forever = Source.repeat("x")Распространённые конструкторы Sink
Аналогично, среди Sink есть Sink.head (первый элемент в виде Future), Sink.seq (собрать всё в Seq), Sink.ignore (вычитать и отбросить) и Sink.last.
Для конвейеров, передающих результат обратно в Ваш код, чаще всего используются Sink.seq и Sink.fold.
import scala.concurrent.Future
val collect: Sink[Int, Future[Seq[Int]]] = Sink.seq
val firstOne: Sink[Int, Future[Int]] = Sink.head
val drain: Sink[Int, Future[akka.Done]] = Sink.ignoreБыстрая проверка
Проверьте, насколько хорошо Вы понимаете линейные типы этапов.
Итоги
Вы изучили три линейных строительных блока: Source создаёт данные, Flow преобразует их, а Sink потребляет. Это неизменяемые, повторно используемые схемы, соединяемые с помощью via и to.
Соединение Source с Sink создаёт RunnableGraph, содержащий материализуемое значение, но не перемещающий данные до вызова run. Далее Вы узнаете, как преобразовывать потоки с помощью более мощных операторов.
Часто задаваемые вопросы
Урок «Источник, поток и приёмник» бесплатный?
Да — полный текст урока «Источник, поток и приёмник» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс Scala for Backend Engineering & Functional Programming, подпишись на CoddyKit PRO. Курс Scala for Backend Engineering & Functional Programming содержит 4 уроков всего.
Чему я научусь в уроке «Источник, поток и приёмник»?
Основные строительные блоки потоковой обработки. Ты практикуешь Scala for Backend Engineering & Functional Programming с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.
Нужен ли мне опыт, чтобы начать Scala for Backend Engineering & Functional Programming?
Предыдущий опыт не требуется. Scala for Backend Engineering & Functional Programming на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 1 из 4.
Сколько времени занимает урок «Источник, поток и приёмник»?
Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.
Можно ли писать и запускать код в этом уроке Scala for Backend Engineering & Functional Programming?
Да. Каждый урок Scala for Backend Engineering & Functional Programming включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.
Все уроки этого курса
- Источник, поток и приёмник
- Преобразование потоков
- Обратное давление
- Запуск конвейера