Kotlin: Flow do Kotlin Explicado

Última atualização: 2026-08-26

Flow é a extensão de streaming das coroutines — streams frios emitem sob demanda, backpressure é embutido, e operadores são composíveis. Charlie usa Flow para construir um pipeline de processamento de pedidos em tempo real.

1. O que Você Aprenderá


2. A História Real de um Arquiteto

(1) Ponto de Dor: Gargalo no Processamento de Pedidos em Tempo Real

O OrderProcessor de Charlie precisa processar 5.000 eventos de pedidos por segundo em tempo real. Processamento em lote com List tem alta latência, aninhamento de callbacks é impossível de manter, e RxJava tem uma curva de aprendizado íngreme.

(2) A Solução com Flow

KOTLIN
// Flow: stream frio, emissão sob demanda, backpressure embutido
fun orderEventStream(): Flow<OrderEvent> = flow {
    while (true) {
        val event = pollOrderEvent()  // Poll não-bloqueante
        emit(event)                   // Emitir para downstream
    }
}

// Compor pipeline: filter -> transform -> collect
orderEventStream()
    .filter { it.total > 100 }
    .map { enrich(it) }
    .collect { dispatch(it) }

Flow é streaming nativo de coroutines — emissão fria sob demanda, backpressure embutido, operadores composíveis, curva de aprendizado muito menor que RxJava.


3. Fundamentos do Flow

(1) Criando Flows

KOTLIN
// 1. construtor flow
fun numberFlow(): Flow<Int> = flow {
    for (i in 1..5) {
        emit(i)  // Emitir valor para o coletor
        delay(100)  // Simular trabalho
    }
}

// 2. flowOf: a partir de valores fixos
val orderFlow = flowOf(
    Order("ORD-001", 299.99),
    Order("ORD-002", 1_500.00)
)

// 3. asFlow: a partir de Iterable/Sequence
val orderList = listOf(Order("ORD-001", 299.99))
val orderFlow2 = orderList.asFlow()

// 4. channelFlow: para emissão concorrente complexa
fun concurrentFlow(): Flow<Int> = channelFlow {
    launch { send(1) }
    launch { send(2) }
}

(2) Comportamento de Stream Frio

KOTLIN
// Flow frio: cada coletor dispara nova emissão
val coldFlow = flow {
    println("Emitindo...")  // Isso executa CADA VEZ que collect é chamado
    emit(1)
    emit(2)
}

coldFlow.collect { println("Coletor 1: $it") }  // Emite 1, 2
coldFlow.collect { println("Coletor 2: $it") }  // Emite 1, 2 NOVAMENTE

(3) Flow Frio vs Flow Quente

Dimensão Flow Frio (Flow) Flow Quente (SharedFlow/StateFlow)
Gatilho Emite apenas quando coletado Executa independentemente dos coletores
Dados Cada coletor obtém dados completos Novos coletores só veem dados subsequentes
Cache Nenhum Replay configurável
Ciclo de vida Segue o coletor Independente
Analogia Vídeo sob demanda Transmissão ao vivo

4. Operações Intermediárias

Operações intermediárias são frias — não disparam emissão, apenas definem lógica de transformação.

(1) map / filter / transform

KOTLIN
// map: transformar cada valor
val totals = orderFlow.map { it.total }

// filter: manter valores correspondentes
val highValue = orderFlow.filter { it.total > 1_000 }

// transform: mais flexível que map (pode emitir múltiplos valores)
val expanded = orderFlow.transform { order ->
    emit("Pedido: ${order.id}")
    emit("Total: \$${order.total} USD")
}

(2) take / drop / debounce

KOTLIN
// take: pegar primeiros N valores
val first3 = orderFlow.take(3)

// drop: pular primeiros N valores
val after5 = orderFlow.drop(5)

// debounce: emitir apenas após período de silêncio
val debounced = searchFlow.debounce(300)  // Esperar 300ms sem novos valores

(3) Categorias de Operações Intermediárias

Categoria Operadores Função
Transformar map / transform Transformação por valor
Filtrar filter / filterNot / filterIsInstance Filtragem condicional
Achatar flatMapConcat / flatMapMerge / flatMapLatest Achatar stream aninhado
Truncar take / drop / takeWhile Pegar/pular primeiros N
Timing debounce / sample / throttleFirst Controle baseado em tempo
Combinar zip / combine Combinação de múltiplos streams

5. Operações Terminais

Operações terminais disparam a emissão e coleta reais.

(1) collect / toList / first

KOTLIN
// collect: consumir cada valor
orderFlow.collect { order -> println(order.id) }

// toList: coletar para lista
val allOrders = orderFlow.toList()

// first: obter primeiro valor
val firstOrder = orderFlow.first()

// single: espera exatamente um valor
val onlyOrder = orderFlow.single()

(2) Categorias de Operações Terminais

Operação Tipo de Retorno Descrição
collect Unit Consumir cada valor
toList / toSet List / Set Coletar para coleção
first / firstOrNull T / T? Obter primeiro valor
single / singleOrNull T / T? Obter valor único
count Int Contar elementos
fold / reduce R / T Acumular

6. Tratamento de Backpressure

(1) O Problema de Backpressure

KOTLIN
// Produtor é mais rápido que o consumidor
val fastFlow = flow {
    repeat(10) {
        emit(it)           // Rápido: 1ms por emissão
        delay(1)
    }
}

// Consumidor lento causa backpressure
fastFlow.collect {
    delay(100)            // Lento: 100ms por processamento
    println("Processado: $it")
}
// Por padrão: produtor aguarda o consumidor (suspend)

(2) Estratégias de Backpressure

KOTLIN
// buffer: executar produtor e consumidor em paralelo
fastFlow.buffer(10)  // Buffer de até 10 valores
    .collect { delay(100); println("Com buffer: $it") }

// conflate: pular valores intermediários, manter o mais recente
fastFlow.conflate()
    .collect { delay(100); println("Mais recente: $it") }

// collectLatest: cancelar coleta anterior quando novo valor chega
fastFlow.collectLatest { value ->
    delay(100)
    println("Processado mais recente: $value")
}

(3) Comparação de Estratégias de Backpressure

Estratégia Comportamento Perda de Dados Caso de Uso
Padrão Produtor aguarda consumidor Sem perda Garantir completude
buffer Armazena N valores em buffer Pausa quando buffer cheio Picos curtos de produção
conflate Pular intermediários, manter mais recente Valores intermediários perdidos Exibição em tempo real
collectLatest Cancelar processamento antigo, tratar mais recente Processamento antigo interrompido Sugestões de busca

7. Pipeline de Stream Frio do Flow

100%
flowchart LR
    A[Fonte<br/>flow emit] --> B[filter<br/>total > 100]
    B --> C[map<br/>enrich]
    C --> D[debounce<br/>300ms]
    D --> E[collect<br/>dispatch]

8. Exemplo Completo: OrderProcessor Stream de Eventos

KOTLIN
// ============================================
// OrderProcessor - Stream de Eventos com Flow
// Recurso: Processamento de eventos de pedidos em tempo real
// ============================================

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

data class Order(val id: String, val total: Double, val status: String, val customer: String)

sealed class OrderEvent {
    data class Created(val order: Order) : OrderEvent()
    data class Updated(val orderId: String, val newStatus: String) : OrderEvent()
    data class Cancelled(val orderId: String, val reason: String) : OrderEvent()
}

// Simular fonte de eventos
fun orderEventStream(): Flow<OrderEvent> = flow {
    val events = listOf(
        OrderEvent.Created(Order("ORD-001", 299.99, "PENDING", "Alice")),
        OrderEvent.Created(Order("ORD-002", 15_000.00, "PENDING", "Bob")),
        OrderEvent.Updated("ORD-001", "CONFIRMED"),
        OrderEvent.Created(Order("ORD-003", 2_500.00, "PENDING", "Charlie")),
        OrderEvent.Cancelled("ORD-003", "Sem estoque"),
        OrderEvent.Updated("ORD-002", "CONFIRMED"),
        OrderEvent.Created(Order("ORD-004", 8_900.00, "PENDING", "Bob")),
        OrderEvent.Updated("ORD-002", "SHIPPED")
    )
    events.forEach { event ->
        emit(event)
        delay(100)
    }
}

// Enriquecer pedidos criados
fun Order.enrich() = copy(status = "ENRICHED")

fun main() = runBlocking {
    println("=== Stream de Eventos de Pedidos ===")

    // Pipeline: filtrar criados -> enriquecer -> coletar
    orderEventStream()
        .onStart { println("Stream iniciada") }
        .onEach { println("  Recebido: $it") }
        .filter { it is OrderEvent.Created }
        .map { (it as OrderEvent.Created).order }
        .filter { it.total > 1_000 }
        .map { it.enrich() }
        .onEach { println("  >>> Enriquecido de alto valor: ${it.id} (\$${it.total} USD)") }
        .onCompletion { println("Stream concluída") }
        .collect()

    // Pipeline separado: atualizações de status
    println("\n=== Atualizações de Status ===")
    orderEventStream()
        .filter { it is OrderEvent.Updated }
        .map { it as OrderEvent.Updated }
        .collect { event ->
            println("  ${event.orderId} -> ${event.newStatus}")
        }

    // Agregação: contar eventos por tipo
    println("\n=== Contagem de Eventos ===")
    val counts = orderEventStream()
        .groupBy { it::class.simpleName ?: "Unknown" }
        .map { (type, flow) ->
            type to flow.count()
        }
        .toList()
        .associate { it }

    counts.forEach { (type, count) ->
        println("  $type: $count eventos")
    }
}

Saída:

TEXT 📖 Somente leitura
=== Stream de Eventos de Pedidos ===
Stream iniciada
  Recebido: Created(order=Order(ORD-001, 299.99, PENDING, Alice))
  Recebido: Created(order=Order(ORD-002, 15000.0, PENDING, Bob))
  >>> Enriquecido de alto valor: ORD-002 ($15000.0 USD)
  Recebido: Updated(orderId=ORD-001, newStatus=CONFIRMED)
  Recebido: Created(order=Order(ORD-003, 2500.0, PENDING, Charlie))
  >>> Enriquecido de alto valor: ORD-003 ($2500.0 USD)
  Recebido: Cancelled(orderId=ORD-003, reason=Sem estoque)
  Recebido: Updated(orderId=ORD-002, newStatus=CONFIRMED)
  Recebido: Created(order=Order(ORD-004, 8900.0, PENDING, Bob))
  >>> Enriquecido de alto valor: ORD-004 ($8900.0 USD)
  Recebido: Updated(orderId=ORD-002, newStatus=SHIPPED)
Stream concluída

=== Atualizações de Status ===
  ORD-001 -> CONFIRMED
  ORD-002 -> CONFIRMED
  ORD-002 -> SHIPPED

=== Contagem de Eventos ===
  Created: 4 eventos
  Updated: 3 eventos
  Cancelled: 1 eventos

9. Exemplos práticos rápidos

▶ Exemplo: Flow básico

KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

fun main() = runBlocking {
    // Flow é um stream frio - valores emitidos sob demanda
    val numbers = flow {
        for (i in 1..3) {
            delay(100)
            emit(i)
        }
    }

    numbers.collect { println("Valor: $it") }
}

// Saída:
// Valor: 1
// Valor: 2
// Valor: 3

▶ Exemplo: Operadores intermediários

KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

fun main() = runBlocking {
    val result = (1..10).asFlow()
        .filter { it % 2 == 0 }
        .map { it * it }
        .take(3)
        .toList()

    println("Quadrados pares: $result")

    // Com reduce
    val sum = (1..5).asFlow().reduce { a, b -> a + b }
    println("Soma: $sum")
}

▶ Exemplo: collect e transformação

KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

fun main() = runBlocking {
    val flow = (1..5).asFlow()
        .onEach { println("Emitindo $it") }
        .map { it * 10 }

    flow.collect { value -> println("Recebido: $value") }
}

// Saída:
// Emitindo 1
// Recebido: 10
// Emitindo 2
// Recebido: 20
// ...

▶ Exemplo: StateFlow e SharedFlow

KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

// StateFlow: emite valor atual para novos assinantes
class OrderCounter {
    private val _count = MutableStateFlow(0)
    val count: StateFlow<Int> = _count.asStateFlow()

    fun increment() {
        _count.value++
    }
}

fun main() = runBlocking {
    val counter = OrderCounter()

    // Assinante 1 - recebe o valor atual (0) e futuros incrementos
    val job1 = launch {
        counter.count.collect { println("Assinante 1: $it") }
    }

    delay(100)
    counter.increment()
    delay(100)
    counter.increment()
    counter.increment()
    delay(100)
    job1.cancel()
}

▶ Exemplo: Combinação de Flows

KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

fun main() = runBlocking {
    val flow1 = (1..5).asFlow().onEach { delay(100) }
    val flow2 = ("A".."E").asFlow().onEach { delay(80) }

    // Combinar valores correspondentes
    flow1.zip(flow2) { n, c -> "$n$c" }
        .collect { println("Combinado: $it") }

    // Concatenar
    val merged = flow {
        emitAll((1..3).asFlow())
        emitAll((4..6).asFlow())
    }
    val collected = merged.toList()
    println("Concatenado: $collected")
}

▶ Exemplo: catch e error handling

KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

fun main() = runBlocking {
    // Sem tratamento: erro se propaga
    val problematicFlow = flow {
        emit(1)
        throw RuntimeException("Erro mid-flow!")
    }

    // Com catch
    problematicFlow
        .catch { e -> println("Capturado: ${e.message}") }
        .collect { println("Valor: $it") }

    // try-catch no collect
    println("\n--- Coletando com try-catch ---")
    val safeFlow = flow {
        listOf("OK", "FAIL", "OK").forEach {
            if (it == "FAIL") error("Falha!")
            emit(it)
        }
    }

    try {
        safeFlow.collect { println("Resultado: $it") }
    } catch (e: Exception) {
        println("Exceção: ${e.message}")
    }
}

▶ Exemplo: Flow com StateIn

KOTLIN
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

fun main() = runBlocking {
    // stateIn: converte Flow frio em StateFlow quente
    val hotFlow = flow {
        for (i in 1..5) {
            emit(i)
            delay(50)
        }
    }.stateIn(
        scope = this,
        started = SharingStarted.WhileSubscribed(500),
        initialValue = 0
    )

    // Múltiplos assinantes compartilham o mesmo stream
    val collector1 = launch {
        hotFlow.collect { println("[1] $it") }
    }
    val collector2 = launch {
        hotFlow.collect { println("[2] $it") }
    }

    delay(500)
    collector1.cancel()
    collector2.cancel()
}

❓ Perguntas Frequentes

P: Qual a diferença entre Flow e Sequence? R: Flow é assíncrono (suporta suspend), Sequence é síncrono. Flow é adequado para fontes de dados assíncronas (rede/banco de dados); Sequence é adequado para processamento síncrono de grandes volumes de dados.

P: Qual a diferença entre Flow e Observable do RxJava? R: Flow é nativo de coroutines, stream-frio-primeiro, com curva de aprendizado suave. RxJava é mais maduro, porém mais complexo. Escolha Flow para novos projetos; coexista com RxJava em projetos existentes.

P: Quando um flow frio se torna um flow quente? R: Use shareIn para converter um Flow frio em SharedFlow, ou stateIn para StateFlow. Isso permite que múltiplos coletores compartilhem a mesma fonte de dados.

P: Como escolher entre flatMapConcat, flatMapMerge, flatMapLatest? R: flatMapConcat processa sequencialmente (padrão), flatMapMerge processa concorrentemente (maior throughput), flatMapLatest cancela streams antigos (cenários em tempo real como busca).

P: Como o Flow trata backpressure por padrão? R: O comportamento padrão é "produtor aguarda consumidor" — emit é uma função suspend; quando o consumidor é lento, o produtor automaticamente suspende. Este é o padrão mais seguro.

P: collect é uma função suspend? R: Sim. collect é uma função suspend, chamável apenas de coroutines ou funções suspend. Ela suspende até o Flow completar.


📖 Resumo


📝 Exercícios

  1. Iniciante (⭐): Crie um Flow que emita 1-10 usando flow { }, depois filtre e imprima números pares. Dica: flow { for (i in 1..10) emit(i) }.filter { }
  2. Intermediário (⭐⭐): Simule um stream de eventos de pedidos com Flow, use filter + map + collect para processar pedidos com valores > 1000. Dica: consulte a Seção 8
  3. Avançado (⭐⭐⭐): Implemente um stream de busca com debounce: simule entrada do usuário, dispare busca apenas após 300ms sem nova entrada. Dica: debounce(300)

← Anterior | Próximo →

Web-Tutorial.com

Equipe Técnica Web-Tutorial

Uma plataforma de tutoriais mantida por diversos desenvolvedores. Cada tutorial é escrito e revisado por profissionais da área correspondente. Trabalhamos para manter nosso conteúdo preciso e confiável — se encontrar algum problema, avise-nos.

100%