Source, Flow y Sink
Los componentes básicos del streaming
Source, Flow y Sink es una lección gratuita de Scala for Backend Engineering & Functional Programming en CoddyKit. Esta es la lección 1 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.
Los tres bloques de construcción
Akka Streams modela una canalización de datos como un grafo de etapas de procesamiento. Las tres etapas lineales principales son Source (produce elementos), Flow (los transforma) y Sink (los consume).
Un Source tiene una salida, un Sink tiene una entrada y un Flow tiene exactamente una entrada y una salida. Conectarlos describe qué debe ocurrir, no cuándo.
Definir un Source
Un Source[Out, Mat] emite elementos de tipo Out y expone un valor materializado de tipo Mat. Los Source más sencillos proceden de colecciones en memoria o rangos.
Hasta que se ejecuta el flujo, un Source es solo un esquema inmutable que puede reutilizarse libremente.
import akka.stream.scaladsl.Source
val numbers: Source[Int, akka.NotUsed] =
Source(1 to 100)
val single: Source[String, akka.NotUsed] =
Source.single("hello")Definir un Sink
Un Sink[In, Mat] consume elementos de tipo In. El valor materializado suele capturar el resultado del consumo, como un Future que se completa cuando termina el flujo.
Sink.foreach ejecuta un efecto secundario por elemento; Sink.fold acumula un único resultado.
import akka.stream.scaladsl.Sink
import scala.concurrent.Future
val printSink: Sink[Int, Future[akka.Done]] =
Sink.foreach(println)
val sumSink: Sink[Int, Future[Int]] =
Sink.fold(0)(_ + _)Definir un Flow
Un Flow[In, Out, Mat] se sitúa entre un Source y un Sink y transforma cada elemento. Los Flow son reutilizables por sí solos y pueden componerse antes de conectarlos a cualquier extremo.
Aquí, un Flow duplica los enteros y los convierte en cadenas.
import akka.stream.scaladsl.Flow
val doubleToString: Flow[Int, String, akka.NotUsed] =
Flow[Int]
.map(_ * 2)
.map(n => s"value=$n")Conectar Source con Sink
El operador via conecta un Flow a un Source, mientras que to conecta un Sink. Conectar directamente un Source con un Sink mediante to produce un RunnableGraph: un esquema cerrado y ejecutable.
Aún no se mueve ningún elemento; esto solo describe la topología.
import akka.stream.scaladsl.{Source, Sink, RunnableGraph}
val graph: RunnableGraph[akka.NotUsed] =
Source(1 to 10).to(Sink.foreach(println))via: insertar un Flow
Use via para insertar un Flow en la canalización. Un Source.via(flow) produce un Source nuevo cuyo tipo de salida coincide con la salida del Flow.
Encadenar llamadas a via permite construir canalizaciones largas de transformación a partir de pequeñas piezas de Flow fáciles de probar.
val pipeline =
Source(1 to 10)
.via(Flow[Int].filter(_ % 2 == 0))
.via(Flow[Int].map(_ * 10))
.to(Sink.foreach(println))Seguridad de tipos entre etapas
El compilador comprueba que el tipo de salida de cada etapa coincida con el tipo de entrada de la siguiente. Un Source[Int] no puede conectarse a un Sink[String] sin un Flow intermedio que convierta el tipo.
Esta comprobación estática detecta errores de conexión de la canalización antes del tiempo de ejecución.
// Source[Int] -> Flow[Int, String] -> Sink[String]
val ok =
Source(1 to 3)
.via(Flow[Int].map(_.toString))
.to(Sink.foreach[String](println))Valores materializados
Cada esquema contiene un valor materializado: un identificador producido cuando se ejecuta el flujo. Un Sink.fold materializa un Future del resultado. De forma predeterminada, combinar etapas conserva el valor materializado situado más a la izquierda (NotUsed para los Source simples).
Use toMat y Keep para seleccionar el valor del lado que desea.
import akka.stream.scaladsl.Keep
import scala.concurrent.Future
val g: RunnableGraph[Future[Int]] =
Source(1 to 100)
.toMat(Sink.fold(0)(_ + _))(Keep.right)Componentes reutilizables
Como los Source, Flow y Sink son valores inmutables, puede definirlos una vez y reutilizarlos en muchas canalizaciones. Esto fomenta una biblioteca de pequeñas etapas de procesamiento con nombre.
Un Flow definido para analizar datos puede insertarse tanto en una canalización de archivos como en una canalización HTTP.
val parse: Flow[String, Int, akka.NotUsed] =
Flow[String].map(_.trim.toInt)
val fromFile = lines.via(parse)
val fromHttp = requestBody.via(parse)Constructores habituales de Source
Akka Streams incluye muchas factorías de Source: Source.single, Source.repeat, Source.tick para emisiones temporizadas, Source.future a partir de un Future y Source.empty.
Elegir el constructor adecuado hace explícita la intención del productor de datos.
import scala.concurrent.duration._
val ticks = Source.tick(0.seconds, 1.second, "tick")
val onceF = Source.future(scala.concurrent.Future.successful(42))
val forever = Source.repeat("x")Constructores habituales de Sink
Del mismo modo, entre los Sink se incluyen Sink.head (el primer elemento como Future), Sink.seq (recopila todos los elementos en un Seq), Sink.ignore (drena y descarta) y Sink.last.
En las canalizaciones que devuelven un resultado al código, Sink.seq y Sink.fold son los más habituales.
import scala.concurrent.Future
val collect: Sink[Int, Future[Seq[Int]]] = Sink.seq
val firstOne: Sink[Int, Future[Int]] = Sink.head
val drain: Sink[Int, Future[akka.Done]] = Sink.ignoreComprobación rápida
Compruebe su comprensión de los tipos de etapas lineales.
Resumen
Ha aprendido los tres bloques de construcción lineales: Source produce, Flow transforma y Sink consume. Son esquemas inmutables y reutilizables que se conectan con via y to.
Conectar un Source con un Sink produce un RunnableGraph que contiene un valor materializado, pero no mueve ningún dato hasta que se ejecuta. A continuación transformará flujos con operadores más completos.
Preguntas frecuentes
¿La lección «Source, Flow y Sink» es gratis?
Sí — el texto completo de «Source, Flow y Sink» 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 «Source, Flow y Sink»?
Los componentes básicos del streaming 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 1 de 4.
¿Cuánto tiempo toma la lección «Source, Flow y Sink»?
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
- Source, Flow y Sink
- Transformar streams
- Backpressure
- Ejecutar un pipeline