Padrões de concorrência
Esta aula aborda os principais padrões de concorrência em Go: worker pool, fan-in/fan-out, pipeline e rate limiting. Cada padrão é explicado com exemplos práticos de código, destacando o uso de goroutines e canais para construir programas concorrentes eficientes e controlados.
Nesta aula, exploraremos quatro padrões fundamentais de concorrência em Go: worker pool, fan-in/fan-out, pipeline e rate limiting. Esses padrões são amplamente utilizados para gerenciar goroutines e canais de forma eficiente, permitindo construir sistemas concorrentes robustos e escaláveis. Cada padrão resolve um problema específico de coordenação e comunicação entre tarefas concorrentes.
Compreender esses padrões é essencial para qualquer desenvolvedor Go que deseje escrever código concorrente idiomático e de alto desempenho. Vamos mergulhar em cada um deles com exemplos práticos.
Worker pool
O padrão worker pool cria um número fixo de goroutines (workers) que processam tarefas de uma fila compartilhada. Isso limita o número de goroutines ativas simultaneamente, evitando sobrecarga e controlando o uso de recursos. É ideal para processar lotes de tarefas independentes, como requisições HTTP ou jobs em lote.
No exemplo abaixo, criamos um pool de 3 workers que leem tarefas de um canal de entrada e enviam resultados para um canal de saída. O canal de tarefas é fechado após todas as tarefas serem enviadas, sinalizando aos workers para terminarem.
package main
import (
"fmt"
"time"
)
func worker(id int, jobs <-chan int, results chan<- int) {
for j := range jobs {
fmt.Printf("Worker %d processing job %d\n", id, j)
time.Sleep(time.Second) // Simula trabalho
results <- j * 2
}
}
func main() {
const numJobs = 5
const numWorkers = 3
jobs := make(chan int, numJobs)
results := make(chan int, numJobs)
// Inicia workers
for w := 1; w <= numWorkers; w++ {
go worker(w, jobs, results)
}
// Envia jobs
for j := 1; j <= numJobs; j++ {
jobs <- j
}
close(jobs)
// Coleta resultados
for a := 1; a <= numJobs; a++ {
<-results
}
}
Neste código, os workers consomem tarefas até que o canal jobs seja fechado. O canal de resultados é usado para coletar as saídas. Note que usamos canais com buffer para evitar bloqueios desnecessários.
Fan-in/fan-out
Fan-out é o padrão de distribuir tarefas para múltiplas goroutines (como no worker pool), enquanto fan-in é o padrão de combinar múltiplos canais de entrada em um único canal de saída. Juntos, eles permitem paralelizar o processamento e agregar resultados.
No exemplo abaixo, criamos duas goroutines que geram números e uma função fan-in que combina os dois canais em um só, usando uma goroutine adicional para multiplexar.
package main
import (
"fmt"
"sync"
)
func generate(nums ...int) <-chan int {
out := make(chan int)
go func() {
for _, n := range nums {
out <- n
}
close(out)
}()
return out
}
func fanIn(chs ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for _, ch := range chs {
wg.Add(1)
go func(c <-chan int) {
for v := range c {
out <- v
}
wg.Done()
}(ch)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
ch1 := generate(1, 2, 3)
ch2 := generate(4, 5, 6)
out := fanIn(ch1, ch2)
for v := range out {
fmt.Println(v)
}
}
O fan-in usa um sync.WaitGroup para aguardar que todas as goroutines de entrada terminem antes de fechar o canal de saída. Isso garante que todos os valores sejam coletados.
Pipeline
O padrão pipeline organiza o processamento em estágios conectados por canais. Cada estágio é uma goroutine que recebe dados de um canal de entrada, processa e envia para o próximo canal. Isso permite construir fluxos de processamento modulares e concorrentes.
No exemplo abaixo, temos três estágios: geração de números, multiplicação por 2 e impressão.
package main
import "fmt"
func gen(nums ...int) <-chan int {
out := make(chan int)
go func() {
for _, n := range nums {
out <- n
}
close(out)
}()
return out
}
func sq(in <-chan int) <-chan int {
out := make(chan int)
go func() {
for n := range in {
out <- n * n
}
close(out)
}()
return out
}
func main() {
// Pipeline
for n := range sq(sq(gen(1, 2, 3, 4))) {
fmt.Println(n)
}
}
Aqui, gen produz números, sq eleva ao quadrado. Compomos dois sq em série para elevar ao quadrado duas vezes. Cada estágio executa concorrentemente, e o pipeline é eficiente porque a comunicação é feita via canais.
Rate limiting
Rate limiting controla a taxa de execução de operações, evitando sobrecarga de recursos. Em Go, podemos implementar rate limiting usando um canal ticker ou um canal com um buffer que funciona como um token bucket.
O exemplo abaixo limita a uma requisição por segundo usando time.Ticker.
package main
import (
"fmt"
"time"
)
func main() {
requests := make(chan int, 5)
for i := 1; i <= 5; i++ {
requests <- i
}
close(requests)
limiter := time.Tick(200 * time.Millisecond) // 5 requisições por segundo
for req := range requests {
<-limiter // Espera o tick
fmt.Println("Request", req, time.Now())
}
}
Para um controle mais flexível, podemos usar um canal com buffer que armazena tokens, permitindo bursts.
package main
import (
"fmt"
"time"
)
func main() {
burstyLimiter := make(chan time.Time, 3)
// Preenche o canal com tokens iniciais (burst de 3)
for i := 0; i < 3; i++ {
burstyLimiter <- time.Now()
}
// Goroutine que adiciona um token a cada 200ms
go func() {
for t := range time.Tick(200 * time.Millisecond) {
burstyLimiter <- t
}
}()
burstyRequests := make(chan int, 5)
for i := 1; i <= 5; i++ {
burstyRequests <- i
}
close(burstyRequests)
for req := range burstyRequests {
<-burstyLimiter
fmt.Println("Request", req, time.Now())
}
}
No segundo exemplo, o burst de 3 tokens permite processar 3 requisições imediatamente, depois o ritmo cai para 5 por segundo.
Boas práticas e observações finais
Ao usar esses padrões, lembre-se de sempre fechar canais quando não houver mais dados a serem enviados, para evitar deadlocks. Use sync.WaitGroup para coordenar a finalização de goroutines. Para rate limiting, considere bibliotecas como golang.org/x/time/rate para implementações mais robustas. Esses padrões são a base para construir sistemas concorrentes em Go e são amplamente utilizados em projetos reais.
Referências
- Go Blog: Pipelines
- Effective Go: Concurrency
- Package rate (golang.org/x/time/rate)
- Go by Example: Worker Pools
- Go by Example: Rate Limiting
- Medium: Concurrency Patterns in Go
Exercícios
Implemente um worker pool que processa 10 tarefas com 4 workers. Cada tarefa deve dormir por um tempo aleatório entre 0 e 2 segundos e depois imprimir seu ID. Use canais com buffer.
✓ Resposta:package main import ( "fmt" "math/rand" "time" ) func worker(id int, jobs <-chan int, results chan<- int) { for j := range jobs { time.Sleep(time.Duration(rand.Intn(2000)) * time.Millisecond) fmt.Printf("Worker %d finished job %d\n", id, j) results <- j } } func main() { const numJobs = 10 const numWorkers = 4 jobs := make(chan int, numJobs) results := make(chan int, numJobs) for w := 1; w <= numWorkers; w++ { go worker(w, jobs, results) } for j := 1; j <= numJobs; j++ { jobs <- j } close(jobs) for a := 1; a <= numJobs; a++ { <-results } }Crie um pipeline de três estágios: o primeiro gera números de 1 a 10, o segundo eleva ao cubo e o terceiro imprime o resultado. Use canais.
✓ Resposta:package main import "fmt" func generate(nums ...int) <-chan int { out := make(chan int) go func() { for _, n := range nums { out <- n } close(out) }() return out } func cube(in <-chan int) <-chan int { out := make(chan int) go func() { for n := range in { out <- n * n * n } close(out) }() return out } func main() { for v := range cube(generate(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)) { fmt.Println(v) } }Implemente fan-in com três canais de entrada que enviam números. Use WaitGroup para coordenar.
✓ Resposta:package main import ( "fmt" "sync" ) func generate(nums ...int) <-chan int { out := make(chan int) go func() { for _, n := range nums { out <- n } close(out) }() return out } func fanIn(chs ...<-chan int) <-chan int { out := make(chan int) var wg sync.WaitGroup for _, ch := range chs { wg.Add(1) go func(c <-chan int) { for v := range c { out <- v } wg.Done() }(ch) } go func() { wg.Wait() close(out) }() return out } func main() { ch1 := generate(1, 2, 3) ch2 := generate(4, 5, 6) ch3 := generate(7, 8, 9) out := fanIn(ch1, ch2, ch3) for v := range out { fmt.Println(v) } }Escreva um programa que limite requisições a 3 por segundo usando um ticker. Faça 10 requisições e imprima o timestamp de cada.
✓ Resposta:package main import ( "fmt" "time" ) func main() { requests := make(chan int, 10) for i := 1; i <= 10; i++ { requests <- i } close(requests) limiter := time.Tick(time.Second / 3) // 3 requisições por segundo for req := range requests { <-limiter fmt.Println("Request", req, time.Now()) } }Modifique o exercício 4 para permitir um burst inicial de 5 requisições imediatas, depois limite a 2 por segundo.
✓ Resposta:package main import ( "fmt" "time" ) func main() { burstyLimiter := make(chan time.Time, 5) for i := 0; i < 5; i++ { burstyLimiter <- time.Now() } go func() { for t := range time.Tick(500 * time.Millisecond) { // 2 por segundo burstyLimiter <- t } }() requests := make(chan int, 10) for i := 1; i <= 10; i++ { requests <- i } close(requests) for req := range requests { <-burstyLimiter fmt.Println("Request", req, time.Now()) } }