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

Преобразование потоков

Отображайте и фильтруйте поступающие данные.

«Преобразование потоков» — бесплатный урок Scala for Backend Engineering & Functional Programming на CoddyKit. Это урок 2 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения Scala for Backend Engineering & Functional Programming, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс Scala for Backend Engineering & Functional Programming содержит 4 уроков всего.

Операторы как преобразования

Akka Streams предоставляет богатый набор операторов для Source и Flow. Они напоминают API коллекций Scala, но работают асинхронно и соблюдают обратное давление.

Каждый оператор возвращает новую схему, поэтому преобразования декларативно объединяются ещё до запуска потока.

map и filter

map применяет синхронную функцию к каждому элементу, а filter удаляет элементы, не удовлетворяющие предикату. Это основные инструменты поэлементного преобразования.

Оба оператора сохраняют порядок элементов и передают вниз по потоку сигналы завершения и ошибки.

val flow =
  Flow[Int]
    .filter(_ % 2 == 0)
    .map(n => n * n)

mapConcat для преобразования одного элемента во множество

Если один входной элемент должен порождать несколько выходных, используйте mapConcat. Он принимает функцию, возвращающую итерируемую последовательность, и разворачивает результаты в поток.

Возврат пустой коллекции фактически удаляет элемент.

val explode: Flow[String, String, akka.NotUsed] =
  Flow[String].mapConcat(line => line.split(",").toList)

val words = Source(List("a,b", "c,d,e"))
  .via(explode)

grouped и sliding

grouped(n) объединяет идущие подряд элементы в Seq размером не более n, что удобно для массовой записи в базу данных. sliding(n) выдаёт перекрывающиеся окна.

Пакетная обработка уменьшает накладные расходы на отдельный элемент в конвейерах с интенсивным вводом-выводом.

val batches: Source[Seq[Int], akka.NotUsed] =
  Source(1 to 1000).grouped(100)

val windows =
  Source(1 to 10).sliding(3, step = 1)

scan и fold

scan выдаёт текущее накопленное значение после каждого элемента, формируя изменяющийся поток состояния. fold выдаёт только итоговое накопленное значение после завершения вышестоящего этапа.

Используйте scan для текущих счётчиков, а fold — для итоговой агрегации.

val running =
  Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15

val total =
  Source(1 to 5).fold(0)(_ + _) // 15

mapAsync для асинхронной работы

mapAsync(parallelism) вызывает функцию, возвращающую Future, и выдаёт результаты в исходном порядке, одновременно выполняя до parallelism Future.

Используйте его для асинхронных вызовов, например обращений к базе данных или HTTP-запросов, когда порядок имеет значение.

import scala.concurrent.Future

val enriched =
  Flow[UserId]
    .mapAsync(parallelism = 4)(id => lookup(id))

def lookup(id: UserId): Future[User] = ???

mapAsyncUnordered

mapAsyncUnordered работает подобно mapAsync, но выдаёт каждый результат сразу после его готовности, не учитывая порядок входных элементов.

Это может повысить пропускную способность, если следующий этап не зависит от порядка: медленный Future больше не блокирует быстрые.

val fast =
  Flow[UserId]
    .mapAsyncUnordered(parallelism = 8)(id => lookup(id))

Состояние преобразования с statefulMapConcat

Для поэлементных преобразований, которым требуется изменяемое локальное состояние, statefulMapConcat создаёт новое состояние при каждой материализации и возвращает итерируемую последовательность результатов.

Это безопасный способ хранить счётчики или буферы, не разделяя состояние между запусками потока.

val withIndex: Flow[String, (Int, String), akka.NotUsed] =
  Flow[String].statefulMapConcat { () =>
    var i = 0
    elem => { i += 1; List((i, elem)) }
  }

Операторы, зависящие от времени

Потоки можно преобразовывать не только по содержимому, но и по времени. throttle ограничивает скорость выдачи, groupedWithin объединяет элементы по размеру или истечении времени, а takeWithin ограничивает продолжительность.

Эти операторы незаменимы для ограничения частоты обращений к внешним API.

import scala.concurrent.duration._

val limited =
  Source(1 to 1000)
    .throttle(10, 1.second)
    .groupedWithin(100, 500.millis)

Обработка ошибок в преобразованиях

Исключение, возникшее внутри оператора, по умолчанию приводит к сбою всего потока. Стратегия надзора может вместо этого выполнить resume (удалить ошибочный элемент) или restart этапа.

Подключите стратегию с помощью withAttributes у Flow.

import akka.stream.{ActorAttributes, Supervision}

val safe =
  Flow[String].map(_.toInt)
    .withAttributes(
      ActorAttributes.supervisionStrategy(_ => Supervision.Resume))

Объединение Flow

Небольшие Flow объединяются в более крупные с помощью via, образуя единый повторно используемый Flow. Это позволяет сосредоточить каждый этап на одной задаче и независимо тестировать его.

Объединённый Flow имеет входной тип первого этапа и выходной тип последнего.

val parse  = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)

val parseAndSquare: Flow[String, Int, akka.NotUsed] =
  parse.via(square)

Быстрая проверка

Подумайте об асинхронных преобразованиях и гарантиях сохранения порядка.

Итоги

Вы изучили операторы преобразования: поэлементные map/filter, преобразование одного элемента во множество с помощью mapConcat, группировку с помощью grouped, накопление с помощью scan и fold, а также асинхронную обработку через mapAsync.

Вы также рассмотрели преобразования с состоянием, операторы времени, такие как throttle, стратегии надзора для ошибок и объединение Flow с помощью via. Далее: как обратное давление обеспечивает безопасность этих этапов.

Часто задаваемые вопросы

Урок «Преобразование потоков» бесплатный?

Да — полный текст урока «Преобразование потоков» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 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 структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 2 из 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