0Pricing
Scala for Backend Engineering & Functional Programming · Aula

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)(_ + _) // 15

mapAsync 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

  1. Origem, fluxo e destino
  2. Transformando fluxos
  3. Contrapressão
  4. Executando um pipeline
← Voltar para Scala for Backend Engineering & Functional Programming