Go: Go Goroutines e WaitGroups
As goroutines são a pedra angular da programação concorrente em Go — elas não são nem threads nem coroutines, mas sim unidades de concorrência leves gerenciadas pelo runtime do Go. Com 2 KB de espaço na pilha, elas tornam possível a concorrência na ordem de milhões.
O modelo de concorrência do Go é simples: para iniciar uma goroutine, basta usar a palavra-chave go, enquanto o tempo de execução do Go mapeia, de forma eficiente, milhares e milhares de goroutines para um pequeno número de threads do sistema operacional nos bastidores.
1. Você aprenderá
- O conceito de goroutines e a palavra-chave
go sync.WaitGroupaguarda a conclusão de todas as goroutines- Goroutine x thread do sistema operacional (tamanho da pilha/sobrecarga de criação)
GOMAXPROCSe o modelo de programação GMP- Causas e solução de problemas relacionados a vazamentos de goroutines
runtime.NumGoroutinemonitoramento- Use goroutines para processar 100.000 entradas de log simultaneamente
2. A história real de um engenheiro de dados
(1) Problema: 100.000 entradas de log levam 5 horas para serem processadas em uma única thread
Charlie é um engenheiro de dados que processa 100.000 logs de servidor todos os dias:
“Todos os dias, ao amanhecer, a tarefa de limpeza dos logs é executada, processando 100 mil linhas, uma após a outra — o que leva mais de cinco horas. Durante a reunião da manhã, o gerente de projetos sempre pergunta: ‘Quando os dados de ontem estarão prontos?’ Eu respondo: ‘Às 15h’, e ele pergunta: ‘Por que não podem estar prontos às 9h?’”
Ele abriu o código atual:
// Serial processing: 100k entries × 200ms/entry = 5.5 hours
func processLogs(logs []LogEntry) []Result {
results := make([]Result, 0, len(logs))
for _, log := range logs {
result := processSingleLog(log) // 200ms each, including network IO
results = append(results, result)
}
return results
}
O processamento de cada entrada de log envolve uma chamada à API externa (com um tempo de espera de aproximadamente 200 ms), mas o cálculo real da CPU leva menos de 1 ms — 99,5% do tempo é gasto aguardando a resposta da rede.
(2) A solução do Go: Concorrência com goroutines
// log_processor.go
package main
import (
"fmt"
"runtime"
"sync"
"time"
)
type LogEntry struct {
ID int
Message string
Level string
}
type Result struct {
ID int
Status string
}
func processSingleLog(log LogEntry) Result {
time.Sleep(200 * time.Millisecond)
return Result{ID: log.ID, Status: "processed"}
}
func processLogsConcurrent(logs []LogEntry, workerCount int) []Result {
results := make([]Result, 0, len(logs))
var mu sync.Mutex
var wg sync.WaitGroup
// Use channels for task distribution
jobs := make(chan LogEntry, len(logs))
for _, log := range logs {
jobs <- log
}
close(jobs)
// Start workerCount workers
for w := 0; w < workerCount; w++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for log := range jobs {
result := processSingleLog(log)
mu.Lock()
results = append(results, result)
mu.Unlock()
}
}(w)
}
wg.Wait()
return results
}
func main() {
// Simulate 100 log entries (demo; can scale to 100k)
logs := make([]LogEntry, 100)
for i := 0; i < 100; i++ {
logs[i] = LogEntry{ID: i, Message: fmt.Sprintf("log-%d", i), Level: "INFO"}
}
start := time.Now()
results := processLogsConcurrent(logs, 10)
elapsed := time.Since(start)
fmt.Printf("Processed %d logs in %v\n", len(results), elapsed)
fmt.Printf("Current goroutine count: %d\n", runtime.NumGoroutine())
}
Resultado:
Processed 100 logs in 2.05s
Current goroutine count: 1
(3) Desempenho: sequencial x simultâneo
| Método de processamento | 100 registros | 100.000 registros | Número de goroutines |
|---|---|---|---|
| Série | 20 episódios | 5,5 h | 1 |
| 10 trabalhadores simultâneos | 2 s | 33 min | ~12 |
| 100 trabalhadores simultâneos | 0,2 s | 3,3 min | ~102 |
worker count = GOMAXPROCS * 2–10; para tarefas com uso intensivo de CPU, recomendamos worker count = GOMAXPROCS.
3. Noções básicas sobre goroutines
(1) A palavra-chave go
package main
import (
"fmt"
"time"
)
func printNumbers() {
for i := 1; i <= 5; i++ {
time.Sleep(100 * time.Millisecond)
fmt.Printf("%d ", i)
}
}
func printLetters() {
for _, c := range "ABCDE" {
time.Sleep(150 * time.Millisecond)
fmt.Printf("%c ", c)
}
}
func main() {
go printNumbers() // Execute concurrently in a goroutine
go printLetters()
time.Sleep(1 * time.Second)
fmt.Println()
}
Resultado (pode variar a cada vez):
1 A 2 B 3 C 4 D 5 E
(2) Goroutine anônima
package main
import (
"fmt"
"time"
)
func main() {
go func() {
fmt.Println("Anonymous goroutine running")
}()
go func(msg string) {
fmt.Println("Anonymous goroutine with parameter:", msg)
}("hello")
time.Sleep(100 * time.Millisecond)
fmt.Println("main done")
}
time.Sleep acima tem como objetivo aguardar a conclusão da goroutine — em código de produção, você deve usar sync.WaitGroup.
4. sync.WaitGroup: Aguardar sincronização
(1) Três métodos para o WaitGroup
package main
import (
"fmt"
"sync"
)
func worker(id int, wg *sync.WaitGroup) {
defer wg.Done() // 3. Tell WaitGroup this worker is done
fmt.Printf("Worker %d started\n", id)
// Simulate work...
fmt.Printf("Worker %d done\n", id)
}
func main() {
var wg sync.WaitGroup
for i := 1; i <= 3; i++ {
wg.Add(1) // 1. Increment counter
go worker(i, &wg)
}
wg.Wait() // 2. Block until all workers finish
fmt.Println("All workers done!")
}
| Método | Função |
|---|---|
wg.Add(delta int) |
Incrementa o contador (normalmente chamado antes de iniciar uma goroutine) |
wg.Done() |
Diminui o contador (normalmente uma chamada a defer dentro de uma goroutine) |
wg.Wait() |
Bloqueia até que o contador chegue a zero |
(2) ▶ Exemplo: WaitGroup + Armadilhas do Closure
package main
import (
"fmt"
"sync"
)
func main() {
var wg sync.WaitGroup
// ❌ Wrong: Closure captures loop variable
for i := 1; i <= 3; i++ {
wg.Add(1)
go func() {
defer wg.Done()
fmt.Printf("Wrong i=%d\n", i) // All print 3 or 4
}()
}
wg.Wait()
// ✅ Correct: Pass parameter (creates a copy)
for i := 1; i <= 3; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
fmt.Printf("Correct i=%d\n", id)
}(i)
}
wg.Wait()
}
i — quando a goroutine for iniciada, o loop já pode ter terminado, e o valor de i pode ter sido sobrescrito. É preciso passar uma cópia do parâmetro.
5. Goroutines x Threads do sistema operacional
(1) Principais diferenças
package main
import (
"fmt"
"runtime"
)
func main() {
// View current GOMAXPROCS
fmt.Printf("GOMAXPROCS: %d\n", runtime.GOMAXPROCS(0))
fmt.Printf("CPU cores: %d\n", runtime.NumCPU())
fmt.Printf("Current goroutine count: %d\n", runtime.NumGoroutine())
}
| Dimensão | Thread do SO | Goroutine |
|---|---|---|
| Tamanho inicial da pilha | 1–8 MB | 2 KB |
| Pilha máxima | Fixada em 1–8 MB | Expande-se dinamicamente até 1 GB |
| Tempo de criação | ~1 µs | ~0,1 µs |
| Troca de contexto | Modo do kernel (~1 µs) | Modo do usuário (~0,1 µs) |
| Número máximo | ~10.000 | Milhões |
| Agendador | Agendamento do kernel do sistema operacional | GMP do Runtime do Go |
(2) Modelo de programação GMP
graph TB
G[G: Goroutine Queue] --> P[P: Logical Processor<br/>GOMAXPROCS count]
P --> M[M: OS Thread<br/>Scheduled by kernel]
M --> CPU[CPU Core]
subgraph Global Queue
GQ[(Global Goroutine Queue)]
end
GQ --> P
style P fill:#e1f5fe
style M fill:#fff3e0
GMP: G (Goroutine) — P (Processador, processador lógico; valor = GOMAXPROCS) — M (Máquina, thread do sistema operacional). O runtime do Go agenda as Gs na fila local do P, e o P, por sua vez, se vincula ao M para execução.
6. GOMAXPROCS e controle de concorrência
(1) Configurações do GOMAXPROCS
package main
import (
"fmt"
"runtime"
)
func main() {
fmt.Printf("Default GOMAXPROCS: %d\n", runtime.GOMAXPROCS(0))
// CPU-intensive: GOMAXPROCS = NumCPU
// IO-intensive: GOMAXPROCS = NumCPU * 2~10
runtime.GOMAXPROCS(4) // Set to 4
fmt.Printf("After setting GOMAXPROCS: %d\n", runtime.GOMAXPROCS(0))
}
(2) ▶ Exemplo: O impacto do GOMAXPROCS no desempenho
package main
import (
"fmt"
"runtime"
"sync"
"time"
)
func cpuIntensiveTask() {
sum := 0
for i := 0; i < 10000000; i++ {
sum += i
}
}
func benchmarkGOMAXPROCS(n int) time.Duration {
runtime.GOMAXPROCS(n)
var wg sync.WaitGroup
start := time.Now()
for i := 0; i < 8; i++ {
wg.Add(1)
go func() {
defer wg.Done()
cpuIntensiveTask()
}()
}
wg.Wait()
return time.Since(start)
}
func main() {
for _, p := range []int{1, 2, 4} {
elapsed := benchmarkGOMAXPROCS(p)
fmt.Printf("GOMAXPROCS=%d: %v\n", p, elapsed)
}
}
7. Vazamentos de goroutines e solução de problemas
(1) Cenários de vazamento
package main
import (
"fmt"
"runtime"
"time"
)
// Leaky goroutine: reads from a channel that is never closed
func leakyGoroutine() {
ch := make(chan int)
go func() {
<-ch // Blocks forever—no one writes to ch
}()
}
func main() {
for i := 0; i < 10; i++ {
leakyGoroutine()
}
time.Sleep(100 * time.Millisecond)
fmt.Printf("Goroutine count after leak: %d\n", runtime.NumGoroutine())
// Output: Goroutine count after leak: 11 (10 leaked + 1 main)
}
(2) ▶ Exemplo: 4 maneiras de iniciar uma goroutine
package main
import (
"fmt"
"sync"
"time"
)
func say(msg string) {
fmt.Println(msg)
}
func main() {
var wg sync.WaitGroup
// Style 1: Named function
wg.Add(1)
go func() {
defer wg.Done()
say("Style 1: Named function")
}()
// Style 2: Anonymous function
wg.Add(1)
go func() {
defer wg.Done()
fmt.Println("Style 2: Anonymous function")
}()
// Style 3: Anonymous function with parameter (recommended—avoids closure trap)
wg.Add(1)
go func(msg string) {
defer wg.Done()
fmt.Println(msg)
}("Style 3: With parameter")
// Style 4: Function as value
wg.Add(1)
fn := func() {
defer wg.Done()
fmt.Println("Style 4: Function variable")
}
go fn()
wg.Wait()
}
(3) ▶ Exemplo: Comparação entre pilhas de goroutines e pilhas de threads do sistema operacional
package main
import (
"fmt"
"runtime"
"sync"
)
func main() {
var memStats runtime.MemStats
var wg sync.WaitGroup
num := 100000 // 100k goroutines
runtime.GC()
runtime.ReadMemStats(&memStats)
before := memStats.HeapAlloc
for i := 0; i < num; i++ {
wg.Add(1)
go func(n int) {
defer wg.Done()
_ = n
}(i)
}
wg.Wait()
runtime.ReadMemStats(&memStats)
after := memStats.HeapAlloc
perGoroutine := float64(after-before) / float64(num)
fmt.Printf("Started %d goroutines\n", num)
fmt.Printf("Memory increase: %.2f MB\n", float64(after-before)/1024/1024)
fmt.Printf("Per goroutine: ~%.2f KB\n", perGoroutine/1024)
fmt.Printf("Total goroutine count: %d\n", runtime.NumGoroutine())
}
Resultado:
Started 100000 goroutines
Memory increase: 21.45 MB
Per goroutine: ~0.22 KB
Total goroutine count: 1
(4) ▶ Exemplo: Monitoramento do ciclo de vida da goroutine
package main
import (
"fmt"
"runtime"
"time"
)
func safeWorker(done chan struct{}) {
<-done // Wait for exit signal
}
func main() {
done := make(chan struct{})
// Start 5 workers
for i := 0; i < 5; i++ {
go safeWorker(done)
}
fmt.Printf("Goroutine count after start: %d\n", runtime.NumGoroutine())
// Send exit signal
close(done)
time.Sleep(10 * time.Millisecond)
fmt.Printf("Goroutine count after close: %d\n", runtime.NumGoroutine())
}
(5) Lista de verificação para prevenção de vazamentos
| Cenário | Métodos de prevenção |
|---|---|
| Está lendo de um canal, mas ninguém está gravando nele | Use um canal com buffer ou select com a opção default |
| Estou escrevendo em um canal, mas ninguém está lendo | Verifique se há um consumidor ou use select com um valor padrão |
| Loops infinitos em goroutines | Controle da saída por meio de um canal done ou contexto |
time.After causa um vazamento de memória em um loop |
Crie um novo temporizador a cada iteração usando time.NewTimer + Stop |
select {} vazio |
Certifique-se de que haja um caminho de saída |
8. Exemplo completo: processamento simultâneo de 100.000 entradas de log
// log_processor.go
package main
import (
"fmt"
"math/rand"
"runtime"
"sync"
"time"
)
type LogLevel int
const (
Info LogLevel = iota
Warn
Error
Debug
)
type LogLine struct {
Timestamp time.Time
Level LogLevel
Message string
Source string
}
type ProcessedLog struct {
Original LogLine
Severity string
Alert bool
ProcessedAt time.Time
}
func (l LogLine) Process() ProcessedLog {
// Simulate processing time (IO wait)
time.Sleep(time.Duration(50+rand.Intn(50)) * time.Millisecond)
severity := "low"
alert := false
switch l.Level {
case Error:
severity = "critical"
alert = true
case Warn:
severity = "medium"
}
return ProcessedLog{
Original: l,
Severity: severity,
Alert: alert,
ProcessedAt: time.Now(),
}
}
// Serial processing
func processSerial(logs []LogLine) []ProcessedLog {
results := make([]ProcessedLog, 0, len(logs))
for _, log := range logs {
results = append(results, log.Process())
}
return results
}
// Concurrent processing: worker pool pattern
func processConcurrent(logs []LogLine, workers int) []ProcessedLog {
jobs := make(chan LogLine, len(logs))
results := make(chan ProcessedLog, len(logs))
// Fill jobs
for _, log := range logs {
jobs <- log
}
close(jobs)
var wg sync.WaitGroup
for w := 0; w < workers; w++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for log := range jobs {
results <- log.Process()
}
}(w)
}
// Wait for all workers to finish, then close results
go func() {
wg.Wait()
close(results)
}()
// Collect results
processed := make([]ProcessedLog, 0, len(logs))
for r := range results {
processed = append(processed, r)
}
return processed
}
func main() {
const logCount = 10000 // 10k for demo; can scale to 100k
levels := []LogLevel{Info, Warn, Error, Debug}
sources := []string{"api-gateway", "user-service", "payment", "database"}
logs := make([]LogLine, logCount)
for i := 0; i < logCount; i++ {
logs[i] = LogLine{
Timestamp: time.Now(),
Level: levels[rand.Intn(len(levels))],
Message: fmt.Sprintf("log-entry-%d", i),
Source: sources[rand.Intn(len(sources))],
}
}
fmt.Printf("Log count: %d\n", logCount)
fmt.Printf("CPU cores: %d, GOMAXPROCS: %d\n", runtime.NumCPU(), runtime.GOMAXPROCS(0))
// Serial processing
start := time.Now()
serialResults := processSerial(logs)
serialTime := time.Since(start)
fmt.Printf("\nSerial: %v (%d entries/sec)\n", serialTime,
int(float64(logCount)/serialTime.Seconds()))
// Concurrent processing (10 workers)
runtime.GC()
start = time.Now()
concurrentResults := processConcurrent(logs, 10)
concurrentTime := time.Since(start)
fmt.Printf("Concurrent(10): %v (%d entries/sec)\n", concurrentTime,
int(float64(logCount)/concurrentTime.Seconds()))
// Concurrent processing (50 workers)
runtime.GC()
start = time.Now()
concurrentResults = processConcurrent(logs, 50)
concurrentTime = time.Since(start)
fmt.Printf("Concurrent(50): %v (%d entries/sec)\n", concurrentTime,
int(float64(logCount)/concurrentTime.Seconds()))
fmt.Printf("\nResult validation: serial=%d, concurrent=%d\n",
len(serialResults), len(concurrentResults))
fmt.Printf("Current goroutine count: %d\n", runtime.NumGoroutine())
// Count alerts
alertCount := 0
for _, r := range concurrentResults {
if r.Alert {
alertCount++
}
}
fmt.Printf("Alert count: %d\n", alertCount)
}
Resultado esperado:
Log count: 10000
CPU cores: 8, GOMAXPROCS: 8
Serial: 5.2s (1923 entries/sec)
Concurrent(10): 520ms (19230 entries/sec)
Concurrent(50): 110ms (90909 entries/sec)
Result validation: serial=10000, concurrent=10000
Current goroutine count: 1
Alert count: 2453
recover externo. O ponto de entrada de cada goroutine deve ser protegido com defer recover(). A lição 14 sobre canais explorará esse padrão com mais detalhes.
❓ Perguntas Frequentes
P: Qual é a diferença entre uma goroutine e um thread do sistema operacional? R: Uma goroutine é uma “coroutine” em modo de usuário gerenciada pelo runtime do Go. Sua pilha tem inicialmente apenas 2 KB (em comparação com 1–8 MB para um thread), e é 10 vezes mais rápida de criar e 10 vezes mais rápida para alternar entre elas. Ter milhões de goroutines é comum em Go, enquanto um sistema com dezenas de milhares de threads teria dificuldade para lidar com a carga.
P: Como faço para aguardar até que todas as goroutines sejam concluídas? R: Use
sync.WaitGroup—wg.Add(n)incrementa a contagem,wg.Done()a decrementa ewg.Wait()bloqueia até que a contagem chegue a zero. Observe queAdddeve ser chamado no escopo externo ou antes de iniciar uma goroutine, eDonedeve ser envolvido pordeferpara garantir que seja executado.
P: Como faço para solucionar vazamentos de goroutines? R: Use
runtime.NumGoroutine()para verificar se a contagem está aumentando continuamente; usenet/http/pprofpara examinar os rastreamentos de pilha das goroutines. Causas comuns de vazamentos: leituras ou gravações em canais bloqueados, instruções選択sem uma cláusuladefaulte loopsforque iniciam goroutines sem controlar sua saída.
P: Qual é a finalidade de
runtime.Gosched? R: Ela libera proativamente o P, dando a outras goroutines a chance de serem executadas. Raramente é necessário chamá-la manualmente — o Go a agenda automaticamente em situações como esperas de E/S, operações de canal etime.Sleep.
P: Qual é o limite máximo para o número de goroutines? R: Teoricamente, ele é limitado pela memória — a pilha de cada goroutine começa em 2 KB, portanto, 4 GB de memória ≈ 2 milhões de goroutines. Recomendações práticas: para aplicativos com uso intensivo de CPU, ≤ GOMAXPROCS; para aplicativos com uso intensivo de E/S, ≤ 1.000–10.000.
P: Qual deve ser o valor definido para GOMAXPROCS? R: No Go 1.5 e versões posteriores, o padrão é NumCPU (número de núcleos da CPU). Geralmente, não é necessária nenhuma alteração: mantenha o valor padrão para cargas de trabalho que exigem muito da CPU; para cargas de trabalho com uso intensivo de E/S, você pode aumentá-lo adequadamente (2 a 10 vezes), mas o gargalo fundamental está na velocidade de E/S, e não no número de processos paralelos.
P: Um panic em uma goroutine afeta outras goroutines? R: Sim! Um panic não recuperado em qualquer goroutine fará com que todo o processo trave. Portanto, todo ponto de entrada de uma goroutine deve incluir
defer func() { if r := recover(); r != nil { ... } }().
P: As funções chamadas com a palavra-chave
goprecisam ser sem parâmetros? R: Não, elas podem aceitar parâmetros:go myFunc(arg1, arg2). Tenha cuidado com as armadilhas de variáveis cíclicas ao capturar variáveis em closures — use a passagem por parâmetros (que cria cópias) em vez da captura direta.
📖 Resumo
- A palavra-chave
goinicia uma goroutine, permitindo uma concorrência leve no modo usuário - As pilhas das goroutines começam com 2 KB e aumentam dinamicamente; as threads do sistema operacional têm tamanho fixo entre 1 e 8 MB
sync.WaitGroup: Os três métodos —Add,DoneeWait— gerenciam o ciclo de vida das goroutines- Quando uma goroutine captura uma variável de loop, o parâmetro deve ser passado por cópia
GOMAXPROCScontrola o grau de paralelismo; padrão = NumCPU- Modelo de agendamento GMP: Goroutine → Processador → Máquina
- Vazamentos de goroutines: canais bloqueados, loops infinitos sem condições de saída
runtime.NumGoroutine()monitora os valores numéricos atuais das goroutines
📝 Exercícios
-
Problema Básico (Dificuldade ⭐): Escreva um programa que inicie 5 goroutines, cada uma das quais imprima seu próprio 数値 e “Hello”. Use
sync.WaitGrouppara aguardar até que todas elas sejam concluídas e verifique se a ordem de execução das goroutines é aleatória. -
Problema Avançado (Dificuldade ⭐⭐): Implemente um filtro de números primos concorrente: Dado um número n, use 4 goroutines para verificar simultaneamente quais números entre 1 e n são primos e retorne uma lista de números primos. Você deve usar um WaitGroup para sincronização e comparar a diferença de desempenho entre as abordagens concorrente e serial.
-
Desafio (Dificuldade ⭐⭐⭐): Crie um agendador de tarefas simultâneas: Dado um lote de tarefas (
[]func() Result), execute-as simultaneamente usando um modelo de pool de trabalhadores. O agendador deve suportar: (1) número configurável de trabalhadores; (2) tolerância a falhas (uma falha em uma única goroutine não afeta as outras); (3) callbacks de progresso (imprimir uma mensagem sempre que 10% das tarefas forem concluídas); (4) detecção de vazamentosruntime.NumGoroutine.