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

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.ignore

Comprobació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

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