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

Запуск конвейера

Материализуйте и выполните граф.

Урок 4 из 413 шагов

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

От схемы к выполнению

До этого момента конвейер был лишь чистой схемой. Материализация — это процесс, превращающий схему в работающие акторы, которые действительно перемещают данные.

Ничего не происходит, пока Вы явно не запустите граф, благодаря чему Akka Streams можно составлять и повторно использовать.

ActorSystem

Для материализации требуется ActorSystem, предоставляющий потоки и диспетчер для этапов потока. В современной Akka система также выступает неявным материализатором.

Обычно один ActorSystem обслуживает всё приложение и множество параллельных потоков.

import akka.actor.ActorSystem

implicit val system: ActorSystem =
  ActorSystem("data-pipeline")
import system.dispatcher // ExecutionContext

runWith

Самый прямой способ запустить Source — использовать runWith: он подключает Sink и выполняет материализацию одним шагом, возвращая материализуемое значение этого Sink.

В данном случае результатом будет Future[Int], завершающийся суммой после окончания потока.

import akka.stream.scaladsl.{Source, Sink}
import scala.concurrent.Future

val total: Future[Int] =
  Source(1 to 100).runWith(Sink.fold(0)(_ + _))

Запуск RunnableGraph

Если Вы уже создали замкнутый RunnableGraph с помощью to или toMat, вызовите run(), чтобы материализовать его. Возвращаемым значением будет то материализуемое значение, которое сохранил граф.

Так построение конвейера чётко отделяется от его выполнения.

import akka.stream.scaladsl.Keep
import scala.concurrent.Future

val graph =
  Source(1 to 100)
    .toMat(Sink.fold(0)(_ + _))(Keep.right)

val result: Future[Int] = graph.run()

Удобные операторы запуска

Source предоставляет сокращённые варианты: runForeach, runFold и runReduce подключают соответствующий Sink и сразу запускают его.

Они удобны для распространённых конечных операций над Source.

import scala.concurrent.Future

val printed: Future[akka.Done] =
  Source(1 to 10).runForeach(println)

val sum: Future[Int] =
  Source(1 to 10).runFold(0)(_ + _)

Работа с Future результата

Конечные Sink возвращают Future, который завершается при окончании потока или возникновении ошибки. Зарегистрируйте обратные вызовы с помощью onComplete, чтобы реагировать на успех или ошибку.

Используйте диспетчер ActorSystem как неявный ExecutionContext для этих обратных вызовов.

import scala.util.{Success, Failure}

total.onComplete {
  case Success(value) => println(s"Sum = $value")
  case Failure(ex)    => println(s"Failed: ${ex.getMessage}")
}

Практический конвейер

Типичный конвейер обработки данных читает данные из источника, преобразует их с помощью потоков, выполняет асинхронный ввод-вывод с помощью mapAsync, объединяет данные в пакеты с помощью grouped и записывает их в приемник.

Каждый этап невелик, а весь конвейер материализуется одним вызовом run.

val done =
  lineSource
    .map(parse)
    .mapAsync(4)(validate)
    .grouped(500)
    .runWith(bulkWriteSink)

Перезапуск потоков после сбоев

Для повышения отказоустойчивости оберните источник или поток в RestartSource.withBackoff, чтобы временные сбои (например, разрыв соединения) автоматически запускали повторно с экспоненциально увеличивающейся задержкой.

Это позволяет поддерживать длительно работающие конвейеры загрузки данных без ручного управления.

import akka.stream.scaladsl.RestartSource
import akka.stream.RestartSettings
import scala.concurrent.duration._

val resilient = RestartSource.withBackoff(
  RestartSettings(1.second, 30.seconds, 0.2))(() => flakySource)

Корректное завершение с помощью KillSwitch

KillSwitch позволяет внешнему коду корректно остановить работающий поток. Вставьте KillSwitches.single с помощью viaMat и сохраните его материализованное значение, чтобы позднее вызвать shutdown().

Это особенно важно для потоков, работающих в течение длительного времени и требующих остановки при завершении приложения.

import akka.stream.{KillSwitches, KillSwitch}
import akka.stream.scaladsl.Keep

val (switch, done) =
  source
    .viaMat(KillSwitches.single)(Keep.right)
    .toMat(Sink.ignore)(Keep.both)
    .run()
// later: switch.shutdown()

Освобождение ресурсов

При завершении приложения вызовите terminate для ActorSystem, чтобы освободить его потоки. Добавьте вызов terminate после Future завершения потока, чтобы завершение прошло упорядоченно.

Утечка ActorSystem не дает JVM завершить работу и удерживает ресурсы открытыми.

done.onComplete { _ =>
  system.terminate()
}

Повторное использование материализатора

Многократная материализация одной и той же схемы создает независимые работающие потоки, использующие общие ресурсы ActorSystem. Сама схема остается неизменяемой и не имеющей побочных эффектов.

Это позволяет безопасно определить конвейер один раз и запускать его по требованию для каждой поступающей задачи.

val blueprint =
  Source(1 to 5).toMat(Sink.seq)(Keep.right)

val run1 = blueprint.run()
val run2 = blueprint.run() // independent execution

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

Подумайте, что требуется потоку для фактической обработки элементов.

Повторение

Запуск конвейера означает материализацию схемы с помощью ActorSystem через run, runWith или удобные операторы; каждый из них возвращает результат типа Future.

Вы рассмотрели реалистичные многоэтапные конвейеры, автоматические перезапуски с задержкой, корректное завершение с помощью KillSwitch, освобождение ресурсов через system.terminate() и безопасное повторное использование неизменяемой схемы в независимых запусках.

Можно начать бесплатно

Изучай Scala с ИИ-репетитором — бесплатно

Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.

Курсы
39
Уроки
143

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

Урок «Запуск конвейера» бесплатный?

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