0Pricing
Scala for Backend Engineering & Functional Programming · Lección

Ejecutar un pipeline

Materialice y ejecute un grafo

Ejecutar un pipeline es una lección gratuita de Scala for Backend Engineering & Functional Programming en CoddyKit. Esta es la lección 4 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de Scala for Backend Engineering & Functional Programming, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de Scala for Backend Engineering & Functional Programming incluye 4 lecciones en total.

Del esquema a la ejecución

Hasta ahora, la canalización ha sido un esquema puro. La materialización es el proceso que convierte ese esquema en actores en ejecución que realmente mueven los datos.

No ocurre nada hasta que ejecuta explícitamente el grafo, lo que permite que Akka Streams sea componible y reutilizable.

El ActorSystem

La materialización requiere un ActorSystem, que proporciona los hilos y el dispatcher en los que se ejecutan las etapas del flujo. En Akka moderno, el sistema también actúa como materializer implícito.

Normalmente, un ActorSystem sirve a una aplicación completa y a muchos flujos simultáneos.

import akka.actor.ActorSystem

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

runWith

La forma más directa de ejecutar un Source es runWith, que asocia un Sink y materializa todo en un solo paso, devolviendo el valor materializado de ese Sink.

En este caso, el resultado es un Future[Int] que se completa con la suma cuando termina el flujo.

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

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

Ejecutar un RunnableGraph

Si ya ha creado un RunnableGraph cerrado con to o toMat, llame a run() para materializarlo. El valor devuelto es el valor materializado que haya conservado el grafo.

Esto separa claramente la construcción de la canalización de su ejecución.

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 conveniencia

Los Sources ofrecen accesos directos: runForeach, runFold y runReduce asocian el Sink correspondiente y ejecutan el flujo de inmediato.

Son opciones concisas para operaciones terminales habituales sobre un 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)(_ + _)

Trabajo con el Future de resultado

Los Sinks terminales devuelven un Future que se completa cuando el flujo termina o falla. Registre callbacks con onComplete para reaccionar ante un resultado correcto o un error.

Use el dispatcher del ActorSystem como ExecutionContext implícito para estos callbacks.

import scala.util.{Success, Failure}

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

Una canalización realista

Una canalización de datos típica lee de un Source, transforma los datos con Flows, realiza E/S asíncrona con mapAsync, agrupa con grouped y escribe en un Sink.

Cada etapa es pequeña y todo se materializa con un único run.

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

Reinicio de flujos fallidos

Para aumentar la resiliencia, envuelva un Source o un Flow con RestartSource.withBackoff para que los fallos transitorios (como una conexión interrumpida) activen un reinicio automático con espera exponencial.

Esto mantiene activas las canalizaciones de ingesta de larga duración sin supervisión 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)

Apagado ordenado con KillSwitch

Un KillSwitch permite que el código externo detenga un flujo en ejecución de forma limpia. Inserte KillSwitches.single mediante viaMat y conserve su valor materializado para llamar posteriormente a shutdown().

Esto es esencial para los flujos de larga duración que deben detenerse cuando se apaga la aplicación.

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

Liberación de recursos

Cuando la aplicación se cierra, termine el ActorSystem para liberar sus hilos. Encadene la llamada a terminate después del Future de finalización del flujo para que el apagado sea ordenado.

Dejar un ActorSystem sin liberar mantiene activa la JVM y conserva los recursos abiertos.

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

Reutilización del materializer

Materializar el mismo esquema varias veces crea flujos independientes en ejecución que comparten los recursos del ActorSystem. El esquema en sí permanece inmutable y sin efectos secundarios.

Esto permite definir una canalización una sola vez y ejecutarla bajo demanda para cada trabajo entrante.

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

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

Comprobación rápida

Considere qué se necesita para que un flujo procese realmente los elementos.

Resumen

Ejecutar una canalización significa materializar un esquema con un ActorSystem mediante run, runWith u operadores de conveniencia, cada uno de los cuales devuelve un resultado Future.

Ha visto canalizaciones realistas de varias etapas, reinicios automáticos con espera progresiva, apagado ordenado mediante KillSwitch, limpieza de recursos con system.terminate() y reutilización segura de un esquema inmutable en ejecuciones independientes.

Preguntas frecuentes

¿La lección «Ejecutar un pipeline» es gratis?

Sí — el texto completo de «Ejecutar un pipeline» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de Scala for Backend Engineering & Functional Programming, actualiza a CoddyKit PRO. El curso de Scala for Backend Engineering & Functional Programming incluye 4 lecciones en total.

¿Qué aprenderé en «Ejecutar un pipeline»?

Materialice y ejecute un grafo Practicas Scala for Backend Engineering & Functional Programming con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.

¿Necesito experiencia previa para empezar Scala for Backend Engineering & Functional Programming?

No se requiere experiencia previa. Scala for Backend Engineering & Functional Programming en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 4 de 4.

¿Cuánto tiempo toma la lección «Ejecutar un pipeline»?

La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.

¿Puedo escribir y ejecutar código en esta lección de Scala for Backend Engineering & Functional Programming?

Sí. Cada lección de Scala for Backend Engineering & Functional Programming incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.

Todas las lecciones de este curso

  1. Source, Flow y Sink
  2. Transformar streams
  3. Backpressure
  4. Ejecutar un pipeline
← Volver a Scala for Backend Engineering & Functional Programming