Transformando fluxos
Mapeie e filtre dados em fluxo.
Transformando fluxos é uma aula grátis de Scala for Backend Engineering & Functional Programming no CoddyKit. Esta é a aula 2 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.
Operadores como transformações
O Akka Streams fornece um conjunto abrangente de operadores em Source e Flow que refletem a API de coleções do Scala, mas são executados de forma assíncrona e respeitam a contrapressão.
Cada operador retorna um novo esquema, portanto as transformações são compostas de forma declarativa antes mesmo de o fluxo ser executado.
map e filter
map aplica uma função síncrona a cada elemento; filter descarta os elementos que não satisfazem um predicado. Esses operadores são os principais recursos para transformações elemento a elemento.
Ambos preservam a ordem e propagam a conclusão e as falhas para os estágios seguintes.
val flow =
Flow[Int]
.filter(_ % 2 == 0)
.map(n => n * n)mapConcat para um para muitos
Quando uma entrada deve produzir várias saídas, use mapConcat. Ele recebe uma função que retorna um iterável e achata os resultados no fluxo.
Retornar uma coleção vazia efetivamente descarta o elemento.
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 e sliding
grouped(n) agrupa elementos consecutivos em uma Seq com até n itens, o que é útil para gravações em massa no banco de dados. sliding(n) emite janelas sobrepostas.
O agrupamento reduz o custo por elemento em pipelines com muitas operações 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 e fold
scan emite o acumulador em execução após cada elemento, criando um fluxo de estado em evolução. fold emite apenas o valor acumulado final, uma vez que a origem tenha sido concluída.
Use scan para contadores em tempo real e fold para agregações terminais.
val running =
Source(1 to 5).scan(0)(_ + _) // 0,1,3,6,10,15
val total =
Source(1 to 5).fold(0)(_ + _) // 15mapAsync para trabalho assíncrono
mapAsync(parallelism) chama uma função que retorna um Future e emite os resultados na ordem, executando simultaneamente até parallelism futures.
Use-o para chamadas assíncronas, como consultas ao banco de dados ou solicitações HTTP, quando a ordem for 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, mas emite cada resultado assim que ele é concluído, ignorando a ordem da entrada.
Ele pode melhorar o desempenho quando o estágio seguinte não depende da ordem, pois um future lento deixa de bloquear os mais rápidos.
val fast =
Flow[UserId]
.mapAsyncUnordered(parallelism = 8)(id => lookup(id))Transformação com estado usando statefulMapConcat
Para transformações por elemento que precisam de estado local mutável, statefulMapConcat cria um novo estado a cada materialização e retorna um iterável de saídas.
Essa é a forma segura de manter contadores ou buffers sem compartilhar estado entre as execuções do fluxo.
val withIndex: Flow[String, (Int, String), akka.NotUsed] =
Flow[String].statefulMapConcat { () =>
var i = 0
elem => { i += 1; List((i, elem)) }
}Operadores baseados em tempo
Os fluxos podem transformar dados com base no tempo, além do conteúdo. throttle limita a taxa de emissão, groupedWithin agrupa por tamanho ou tempo decorrido, e takeWithin limita a duração.
Eles são essenciais para limitar a taxa de APIs externas.
import scala.concurrent.duration._
val limited =
Source(1 to 1000)
.throttle(10, 1.second)
.groupedWithin(100, 500.millis)Tratamento de erros nas transformações
Uma exceção lançada dentro de um operador, por padrão, faz todo o fluxo falhar. Uma estratégia de supervisão pode, em vez disso, executar resume (descartar o elemento inválido) ou restart o estágio.
Anexe a estratégia com withAttributes ao Flow.
import akka.stream.{ActorAttributes, Supervision}
val safe =
Flow[String].map(_.toInt)
.withAttributes(
ActorAttributes.supervisionStrategy(_ => Supervision.Resume))Compondo Flows
Flows pequenos podem ser compostos em outros maiores com via, produzindo um único Flow reutilizável. Isso mantém cada transformação concentrada e permite testá-la de forma independente.
O Flow composto tem o tipo de entrada do primeiro e o tipo de saída do ú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)Verificação rápida
Considere as transformações assíncronas e suas garantias de ordenação.
Recapitulação
Você explorou operadores de transformação: map/filter elemento a elemento, mapConcat de um para muitos, agrupamento com grouped, acumulação com scan e fold, e trabalho assíncrono por meio de mapAsync.
Você também conheceu transformações com estado, operadores baseados em tempo, como throttle, estratégias de supervisão para erros e a composição de Flows com via. A seguir: como a contrapressão mantém esses estágios seguros.
Perguntas Frequentes
A aula “Transformando fluxos” é grátis?
Sim — o texto completo de “Transformando fluxos” é 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 “Transformando fluxos”?
Mapeie e filtre dados em fluxo. 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 2 de 4.
Quanto tempo leva a aula “Transformando fluxos”?
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
- Origem, fluxo e destino
- Transformando fluxos
- Contrapressão
- Executando um pipeline