0Pricing
Scala for Backend Engineering & Functional Programming · Aula

Executando um pipeline

Materialize e execute um grafo.

Executando um pipeline é uma aula grátis de Scala for Backend Engineering & Functional Programming no CoddyKit. Esta é a aula 4 de 4. Você pode ler a aula completa abaixo gratuitamente — depois pratica ao vivo no navegador com um editor de código integrado e um tutor de IA 24/7. Faz parte do caminho de aprendizado de Scala for Backend Engineering & Functional Programming, e seu progresso é sincronizado entre a web e o app CoddyKit. O curso de Scala for Backend Engineering & Functional Programming inclui 4 aulas no total.

Do esquema à execução

Até agora, o pipeline foi apenas um esquema puro. A materialização é o processo que transforma esse esquema em atores em execução que realmente movimentam os dados.

Nada acontece até que você execute explicitamente o grafo, e é isso que torna o Akka Streams componível e reutilizável.

O ActorSystem

A materialização requer um ActorSystem, que fornece as threads e o despachante que dão suporte aos estágios do fluxo. No Akka moderno, o sistema também atua como o materializador implícito.

Normalmente, um ActorSystem atende a toda uma aplicação e a muitos fluxos simultâneos.

import akka.actor.ActorSystem

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

runWith

A maneira mais direta de executar uma Source é runWith, que conecta um Sink e o materializa em uma única etapa, retornando o valor materializado desse Sink.

Aqui, o resultado é um Future[Int] que é concluído com a soma quando o fluxo termina.

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

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

run em um RunnableGraph

Se você já tiver criado um RunnableGraph fechado com to ou toMat, chame run() para materializá-lo. O valor retornado é qualquer valor materializado que o grafo tenha mantido.

Isso separa claramente a construção do pipeline da execução.

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()

Operadores run de conveniência

As Sources oferecem atalhos: runForeach, runFold e runReduce conectam cada um o Sink correspondente e executam imediatamente.

Eles são concisos para operações terminais comuns em uma 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)(_ + _)

Trabalhando com o Future do resultado

Os Sinks terminais retornam um Future que é concluído quando o fluxo termina ou falha. Registre retornos de chamada com onComplete para reagir ao sucesso ou ao erro.

Use o despachante do ActorSystem como o ExecutionContext implícito para esses retornos de chamada.

import scala.util.{Success, Failure}

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

Um pipeline realista

Um pipeline de dados típico lê de uma Source, transforma os dados com Flows, realiza E/S assíncrona com mapAsync, agrupa com grouped e grava em um Sink.

Cada estágio é pequeno, e todo o conjunto é materializado com um único run.

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

Reiniciando fluxos com falha

Para obter resiliência, envolva uma Source ou um Flow com RestartSource.withBackoff, para que falhas temporárias (como uma conexão interrompida) acionem uma reinicialização automática com espera exponencial.

Isso mantém pipelines de ingestão de longa duração ativos sem supervisão manual.

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)

Encerramento controlado com KillSwitch

Um KillSwitch permite que um código externo pare um fluxo em execução de forma limpa. Insira KillSwitches.single por meio de viaMat e mantenha seu valor materializado para chamar shutdown() mais tarde.

Isso é essencial para fluxos de longa duração que precisam parar quando a aplicação for encerrada.

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()

Liberando recursos

Quando a aplicação for encerrada, termine o ActorSystem para liberar suas threads. Encadeie a chamada de encerramento após o Future de conclusão do fluxo, para que o desligamento ocorra de forma ordenada.

Deixar um ActorSystem aberto mantém a JVM ativa e conserva os recursos ocupados.

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

Reutilizando o materializador

Materializar o mesmo esquema várias vezes cria fluxos independentes em execução que compartilham os recursos do ActorSystem. O próprio esquema permanece imutável e sem efeitos colaterais.

Isso permite definir um pipeline uma vez e executá-lo sob demanda para cada novo trabalho recebido.

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

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

Verificação rápida

Considere o que é necessário para que um fluxo realmente processe elementos.

Recapitulação

Executar um pipeline significa materializar um esquema com um ActorSystem por meio de run, runWith ou operadores de conveniência, cada um retornando um resultado Future.

Você conheceu pipelines realistas com vários estágios, reinicializações automáticas com espera progressiva, encerramento controlado por meio de KillSwitch, liberação de recursos com system.terminate() e reutilização segura de um esquema imutável em execuções independentes.

Perguntas Frequentes

A aula “Executando um pipeline” é grátis?

Sim — o texto completo de “Executando um pipeline” é grátis para ler aqui na web. Para praticá-la interativamente (um editor de código integrado e um tutor de IA 24/7) e desbloquear o restante do curso de Scala for Backend Engineering & Functional Programming, atualize para CoddyKit PRO. O curso de Scala for Backend Engineering & Functional Programming inclui 4 aulas no total.

O que vou aprender em “Executando um pipeline”?

Materialize e execute um grafo. Você pratica Scala for Backend Engineering & Functional Programming com código prático que executa diretamente no navegador, e um tutor de IA 24/7 responde suas dúvidas enquanto trabalha na aula.

Preciso ter experiência prévia para começar Scala for Backend Engineering & Functional Programming?

Nenhuma experiência prévia é necessária. Scala for Backend Engineering & Functional Programming no CoddyKit é estruturado para alunos iniciantes até avançados, então você pode começar aqui ou desde o início e aprender no seu ritmo. Esta é a aula 4 de 4.

Quanto tempo leva a aula “Executando um pipeline”?

A maioria das aulas CoddyKit leva cerca de 5–10 minutos. Cada uma é compacta e interativa, então você faz progresso constante e retoma exatamente de onde parou entre web e app.

Posso escrever e executar código nesta aula de Scala for Backend Engineering & Functional Programming?

Sim. Cada aula de Scala for Backend Engineering & Functional Programming inclui um editor de código integrado, então você escreve e executa código real direto no navegador e recebe feedback de IA instantaneamente — nenhuma configuração local necessária.

Todas as aulas deste curso

  1. Origem, fluxo e destino
  2. Transformando fluxos
  3. Contrapressão
  4. Executando um pipeline
← Voltar para Scala for Backend Engineering & Functional Programming