0Pricing
Scala for Backend Engineering & Functional Programming · Урок

Источник, поток и приёмник

Основные строительные блоки потоковой обработки.

«Источник, поток и приёмник» — бесплатный урок 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 — локальная установка не требуется.

Все уроки этого курса

  1. Источник, поток и приёмник
  2. Преобразование потоков
  3. Обратное давление
  4. Запуск конвейера
← Назад к Scala for Backend Engineering & Functional Programming