Преобразование потоков
Отображайте и фильтруйте поступающие данные.
«Преобразование потоков» — бесплатный урок 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)(_ + _) // 15mapAsync для асинхронной работы
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 — локальная установка не требуется.
Все уроки этого курса
- Источник, поток и приёмник
- Преобразование потоков
- Обратное давление
- Запуск конвейера