Transformar streams
Mapee y filtre datos en movimiento
Transformar streams es una lección gratuita de Scala for Backend Engineering & Functional Programming en CoddyKit. Esta es la lección 2 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 operadores como transformaciones
Akka Streams proporciona un amplio conjunto de operadores para Source y Flow que reflejan la API de colecciones de Scala, pero se ejecutan de forma asíncrona y respetan la contrapresión.
Cada operador devuelve un esquema nuevo, por lo que las transformaciones se componen de forma declarativa antes de que se ejecute el flujo.
map y filter
map aplica una función síncrona a cada elemento; filter descarta los elementos que no cumplen un predicado. Son los operadores fundamentales de la transformación elemento a elemento.
Ambos conservan el orden y propagan la finalización y los fallos hacia abajo.
val flow =
Flow[Int]
.filter(_ % 2 == 0)
.map(n => n * n)mapConcat para transformar uno en varios
Cuando una entrada debe producir varias salidas, use mapConcat. Recibe una función que devuelve un iterable y aplana los resultados en el flujo.
Devolver una colección vacía descarta el elemento de forma efectiva.
val explode: Flow[String, String, akka.NotUsed] =
Flow[String].mapConcat(line => line.split(",").toList)
val words = Source(List("a,b", "c,d,e"))
.via(explode)grouped y sliding
grouped(n) agrupa elementos consecutivos en un Seq de hasta n elementos, lo que resulta útil para escrituras masivas en bases de datos. sliding(n) emite ventanas superpuestas.
La agrupación reduce la sobrecarga por elemento en canalizaciones con un uso intensivo de E/S.
val batches: Source[Seq[Int], akka.NotUsed] =
Source(1 to 1000).grouped(100)
val windows =
Source(1 to 10).sliding(3, step = 1)scan y fold
scan emite el acumulador actualizado después de cada elemento, lo que proporciona un flujo de estado evolutivo. fold emite únicamente el valor acumulado final cuando se completa el flujo ascendente.
Use scan para contadores en tiempo real y fold para agregaciones terminales.
val running =
Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15
val total =
Source(1 to 5).fold(0)(_ + _) // 15mapAsync para trabajo asíncrono
mapAsync(parallelism) llama a una función que devuelve un Future y emite los resultados en orden, ejecutando de forma concurrente hasta parallelism futures.
Úselo para llamadas asíncronas, como consultas a bases de datos o solicitudes HTTP, cuando el orden sea importante.
import scala.concurrent.Future
val enriched =
Flow[UserId]
.mapAsync(parallelism = 4)(id => lookup(id))
def lookup(id: UserId): Future[User] = ???mapAsyncUnordered
mapAsyncUnordered se comporta como mapAsync, pero emite cada resultado en cuanto se completa, ignorando el orden de entrada.
Puede mejorar el rendimiento cuando al flujo descendente no le importa el orden, ya que un future lento deja de bloquear a los más rápidos.
val fast =
Flow[UserId]
.mapAsyncUnordered(parallelism = 8)(id => lookup(id))Transformación con estado mediante statefulMapConcat
Para transformaciones por elemento que necesitan un estado local mutable, statefulMapConcat crea un estado nuevo en cada materialización y devuelve un iterable de resultados.
Es la forma segura de mantener contadores o búferes sin compartir el estado entre ejecuciones del flujo.
val withIndex: Flow[String, (Int, String), akka.NotUsed] =
Flow[String].statefulMapConcat { () =>
var i = 0
elem => { i += 1; List((i, elem)) }
}Operadores basados en el tiempo
Los flujos pueden transformarse según el tiempo, además de según el contenido. throttle limita la velocidad de emisión, groupedWithin agrupa por cantidad o por tiempo transcurrido, y takeWithin limita la duración.
Estos operadores son esenciales para limitar la velocidad de las API externas.
import scala.concurrent.duration._
val limited =
Source(1 to 1000)
.throttle(10, 1.second)
.groupedWithin(100, 500.millis)Gestión de errores en las transformaciones
Una excepción lanzada dentro de un operador hace que todo el flujo falle de forma predeterminada. Una estrategia de supervisión puede hacer que, en su lugar, se use resume (descartar el elemento incorrecto) o restart (reiniciar la etapa).
Asocie la estrategia al Flow mediante withAttributes.
import akka.stream.{ActorAttributes, Supervision}
val safe =
Flow[String].map(_.toInt)
.withAttributes(
ActorAttributes.supervisionStrategy(_ => Supervision.Resume))Composición de Flows
Los Flows pequeños se pueden combinar en otros más grandes mediante via, lo que produce un único Flow reutilizable. Así, cada transformación se mantiene enfocada y se puede probar de forma independiente.
El Flow compuesto tiene el tipo de entrada del primero y el tipo de salida del último.
val parse = Flow[String].map(_.toInt)
val square = Flow[Int].map(n => n * n)
val parseAndSquare: Flow[String, Int, akka.NotUsed] =
parse.via(square)Comprobación rápida
Considere las transformaciones asíncronas y sus garantías de orden.
Resumen
Ha explorado operadores de transformación: map/filter elemento a elemento, mapConcat de uno a muchos, agrupación con grouped, acumulación con scan y fold, y trabajo asíncrono mediante mapAsync.
También ha visto transformaciones con estado, operadores basados en el tiempo como throttle, estrategias de supervisión para errores y la composición de Flows con via. A continuación: cómo la contrapresión mantiene seguras estas etapas.
Preguntas frecuentes
¿La lección «Transformar streams» es gratis?
Sí — el texto completo de «Transformar streams» 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 «Transformar streams»?
Mapee y filtre datos en movimiento 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 2 de 4.
¿Cuánto tiempo toma la lección «Transformar streams»?
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