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á
- Stream frio vs stream quente:
Flowemite sob demanda - Operações intermediárias:
map/filter/transform/take/debounce - Operações terminais:
collect/toList/first/single - Tratamento de backpressure:
buffer/conflate/collectLatest - Charlie em ação: pipeline de streaming de eventos de pedidos
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
// 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
// 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
// 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
// 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
// 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
// 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
// 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
// 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
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
// ============================================
// 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:
=== 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
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
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
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
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
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
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
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
shareInpara converter um Flow frio em SharedFlow, oustateInpara StateFlow. Isso permite que múltiplos coletores compartilhem a mesma fonte de dados.
P: Como escolher entre flatMapConcat, flatMapMerge, flatMapLatest? R:
flatMapConcatprocessa sequencialmente (padrão),flatMapMergeprocessa concorrentemente (maior throughput),flatMapLatestcancela 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
- Flow é um stream frio: cada coletor dispara emissão independente, sem compartilhamento de dados
- Operações intermediárias (map/filter/transform) são frias — não disparam emissão
- Operações terminais (collect/toList/first) disparam o fluxo real de dados
- Backpressure padrão: produtor aguarda consumidor; ajuste com buffer/conflate/collectLatest
- A sintaxe dos operadores do Flow é consistente com operadores de coleção — baixa curva de aprendizado
onStart/onEach/onCompletionsão hooks do ciclo de vida do stream
📝 Exercícios
- 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 { } - Intermediário (⭐⭐): Simule um stream de eventos de pedidos com Flow, use
filter+map+collectpara processar pedidos com valores > 1000. Dica: consulte a Seção 8 - 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)