Go: Padrões de Go Select e Concorrência
選択é o canivete suíço da programação concorrente em Go — ele monitora vários canais simultaneamente e processa aquele que ficar pronto primeiro. Ele funciona como um “comutador inteligente” que conecta goroutines.
Se pensarmos em um canal como uma linha telefônica que conecta goroutines, então o 選択 é a telefonista — ele escuta todas as chamadas ao mesmo tempo e atende aquela que estiver tocando. Nesta lição, você vai dominar todos os padrões essenciais do 選択.
1. Você aprenderá
- Sintaxe e comportamento da multiplexação
選択 - Seleção aleatória (vários casos prontos ao mesmo tempo)
time.AfterControle de tempo limitedefault: operação sem bloqueiofor-選択Saída do loop (canal concluído)- Distribuir tarefas usando o modo “fan-out”
- Resultados da agregação no modo fan-in
- Criar um pipeline de processamento de dados utilizando o modelo de pipeline
2. A história real de um engenheiro de negociação quantitativa
(1) Problema: A consulta a três fontes de dados faz com que a CPU funcione em plena capacidade
Bob é engenheiro de back-end na equipe de negociação quantitativa. Ele precisa monitorar três fontes de dados de ações simultaneamente:
“Nossa estratégia exige a obtenção de cotações em tempo real de três bolsas simultaneamente: a NYSE, a NASDAQ e a LSE. A implementação anterior em Java utilizava um único thread para consultar três WebSockets, com uma espera ativa de 10 milissegundos a cada vez, consumindo 30% da CPU — meu chefe disse que eu estava consumindo mais energia do que o próprio sistema de negociação.”
Ele abriu o código atual:
// Bad code: busy polling
func pollDataSources() {
for {
// Poll every 10ms, wasting CPU
if data1 := pollNYSE(); data1 != nil {
process(data1)
}
if data2 := pollNASDAQ(); data2 != nil {
process(data2)
}
if data3 := pollLSE(); data3 != nil {
process(data3)
}
time.Sleep(10 * time.Millisecond)
}
}
Três problemas: (1) Um intervalo de sondagem de 10 milissegundos desperdiça recursos da CPU; (2) A chegada e o processamento dos dados não estão sincronizados; (3) Os atrasos na sondagem, que variam de 0 a 10 milissegundos, são imprevisíveis.
(2) Solução do Go: use select para ouvir três canais
// market_data.go
package main
import (
"fmt"
"math/rand"
"time"
)
type Quote struct {
Source string
Symbol string
Price float64
}
func simulateExchange(name string, out chan<- Quote) {
symbols := []string{"AAPL", "GOOGL", "MSFT", "AMZN"}
for {
time.Sleep(time.Duration(50+rand.Intn(200)) * time.Millisecond)
quote := Quote{
Source: name,
Symbol: symbols[rand.Intn(len(symbols))],
Price: 100 + rand.Float64()*200,
}
out <- quote
}
}
func main() {
nyse := make(chan Quote)
nasdaq := make(chan Quote)
lse := make(chan Quote)
go simulateExchange("NYSE", nyse)
go simulateExchange("NASDAQ", nasdaq)
go simulateExchange("LSE", lse)
// select listens to all three channels simultaneously
timeout := time.After(2 * time.Second)
for {
select {
case q := <-nyse:
fmt.Printf("[NYSE] %s: $%.2f\n", q.Symbol, q.Price)
case q := <-nasdaq:
fmt.Printf("[NASDAQ] %s: $%.2f\n", q.Symbol, q.Price)
case q := <-lse:
fmt.Printf("[LSE] %s: $%.2f\n", q.Symbol, q.Price)
case <-timeout:
fmt.Println("Demo ended")
return
}
}
}
Resultado:
[NASDAQ] AMZN: $198.32
[NYSE] AAPL: $150.45
[LSE] MSFT: $287.10
...
(3) Desempenho: SELECT Listen x Polling
| Método | Utilização da CPU | Latência de resposta | Complexidade do código |
|---|---|---|---|
| Intervalo de sondagem de 10 ms | 30% | 0–10 ms | Baixo |
| Intervalo de sondagem de 100 ms | 3% | 0–100 ms | Baixo |
| Escutar | 0% | 0 ms (em tempo real) | Moderado |
選択 + channel é orientado a eventos — a goroutine fica em espera (sem consumir CPU) quando não há dados e é ativada pelo ambiente de execução quando os dados chegam. É isso que torna o modelo de concorrência do Go tão eficiente.
3. Noções básicas sobre SELECT
(1) Sintaxe de seleção
select {
case v := <-ch1:
// ch1 is ready
case v := <-ch2:
// ch2 is ready
case ch3 <- value:
// Can send to ch3 (ch3 has space or has a receiver)
default:
// All channels are not ready
}
(2) ▶ Exemplo: selecionar (seleção aleatória)
package main
import (
"fmt"
)
func main() {
ch1 := make(chan string, 1)
ch2 := make(chan string, 1)
ch1 <- "from ch1"
ch2 <- "from ch2"
// Both channels are ready; select picks one at random
for i := 0; i < 2; i++ {
select {
case msg := <-ch1:
fmt.Println(msg)
case msg := <-ch2:
fmt.Println(msg)
}
}
}
Resultado (pode variar a cada vez):
from ch1
from ch2
4. Controle de tempo limite
(1) Modo de tempo limite
package main
import (
"fmt"
"time"
)
func slowOperation() string {
time.Sleep(2 * time.Second)
return "result"
}
func main() {
ch := make(chan string)
go func() {
ch <- slowOperation()
}()
select {
case result := <-ch:
fmt.Println("Success:", result)
case <-time.After(1 * time.Second):
fmt.Println("Timeout! Operation exceeded 1 second")
}
}
Resultado:
Timeout! Operation exceeded 1 second
(2) ▶ Exemplo: Tempos de espera em níveis (espera em níveis)
package main
import (
"fmt"
"time"
)
func fetchFromCache() string {
time.Sleep(50 * time.Millisecond)
return "cache-data"
}
func fetchFromDB() string {
time.Sleep(200 * time.Millisecond)
return "db-data"
}
func fetchFromAPI() string {
time.Sleep(500 * time.Millisecond)
return "api-data"
}
func main() {
cache := make(chan string)
db := make(chan string)
api := make(chan string)
go func() { cache <- fetchFromCache() }()
go func() { db <- fetchFromDB() }()
go func() { api <- fetchFromAPI() }()
select {
case r := <-cache:
fmt.Println("Cache hit:", r)
case <-time.After(100 * time.Millisecond):
select {
case r := <-db:
fmt.Println("DB returned:", r)
case <-time.After(300 * time.Millisecond):
select {
case r := <-api:
fmt.Println("API returned:", r)
case <-time.After(600 * time.Millisecond):
fmt.Println("All data sources timed out!")
}
}
}
}
time.After(d) retorna <-chan time.Time e envia a hora atual após d segundos. Cada vez que select chama time.After, ele cria um novo temporizador — se você chamá-lo com frequência em um loop, certifique-se de usar time.NewTimer e interrompê-lo para evitar vazamentos de recursos.
5. padrão: operação não bloqueante
(1) Envio/recebimento sem bloqueio
package main
import (
"fmt"
)
func main() {
ch := make(chan int, 1)
// Non-blocking receive
select {
case v := <-ch:
fmt.Println("Received:", v)
default:
fmt.Println("No data (non-blocking)")
}
// Non-blocking send
ch <- 1
select {
case ch <- 2:
fmt.Println("Send successful")
default:
fmt.Println("Buffer full (non-blocking)")
}
}
Resultado:
No data (non-blocking)
Buffer full (non-blocking)
(2) ▶ Exemplo: Canal sem bloqueio + disjuntor de polling
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan int, 3)
go func() {
for i := 1; i <= 10; i++ {
select {
case ch <- i:
// Send successful
default:
fmt.Printf("Buffer full, dropped %d\n", i)
}
time.Sleep(10 * time.Millisecond)
}
close(ch)
}()
// Slow consumer
for v := range ch {
fmt.Printf("Processing: %d\n", v)
time.Sleep(50 * time.Millisecond)
}
}
Resultado:
Processing: 1
Processing: 2
Processing: 3
Buffer full, dropped 4
Buffer full, dropped 5
Processing: 6
...
default faz com que select retorne imediatamente — default é executado se nenhum dos canais estiver pronto. Isso é ideal para operações de canal não bloqueantes, mas tenha cuidado: usar default em um loop for pode causar um loop infinito.
6. Ciclos for-select e o canal done
(1) Modo de saída “canal concluído”
package main
import (
"fmt"
"time"
)
func worker(done <-chan struct{}) {
for {
select {
case <-done:
fmt.Println("worker exiting")
return
default:
fmt.Println("worker working...")
time.Sleep(200 * time.Millisecond)
}
}
}
func main() {
done := make(chan struct{})
go worker(done)
time.Sleep(1 * time.Second)
close(done)
time.Sleep(100 * time.Millisecond)
fmt.Println("main exiting")
}
(2) ▶ Exemplo: Três maneiras de sair de um loop “for-selecção”
package main
import (
"fmt"
"time"
)
func main() {
ch := make(chan int)
done := make(chan struct{})
// Producer
go func() {
for i := 1; i <= 5; i++ {
ch <- i
time.Sleep(100 * time.Millisecond)
}
close(ch)
}()
// Consumer (for-select loop)
go func() {
for {
select {
case v, ok := <-ch:
if !ok {
fmt.Println("Style 1: channel closed exit")
close(done)
return
}
fmt.Printf("Processing: %d\n", v)
case <-time.After(1 * time.Second):
fmt.Println("Style 2: timeout exit")
close(done)
return
}
}
}()
<-done
fmt.Println("main exiting")
}
(3) Comparação dos modos de saída do for-select
| Modo | Condições de acionamento | Vantagens | Desvantagens |
|---|---|---|---|
| close(ch) | O remetente fecha o canal | Encerramento natural | Só pode ser usado pelo destinatário |
| canal pronto | sinal close(done) |
Flexível; pode ser acionado externamente | Requer um canal adicional |
| tempo limite | time.After tempo limite |
evita impasses | os tempos limite rígidos não são flexíveis o suficiente |
| contexto | ctx.Done() |
Pode enviar um sinal de cancelamento | Lição 17: Aprofundamento |
7. Modo Fan-out / Fan-in
(1) Fan-out: Distribuição de tarefas
package main
import (
"fmt"
"sync"
)
func fanOut(jobs <-chan int, workers int) []<-chan int {
channels := make([]<-chan int, workers)
for w := 0; w < workers; w++ {
ch := make(chan int, 10)
channels[w] = ch
go func(id int, out chan<- int) {
defer close(out)
for job := range jobs {
result := job * job
fmt.Printf("Worker %d: %d^2 = %d\n", id, job, result)
out <- result
}
}(w, ch)
}
return channels
}
(2) ▶ Exemplo: fan-in: agregação de resultados
package main
import (
"fmt"
"sync"
)
func fanIn(channels ...<-chan int) <-chan int {
out := make(chan int)
var wg sync.WaitGroup
for _, ch := range channels {
wg.Add(1)
go func(c <-chan int) {
defer wg.Done()
for v := range c {
out <- v
}
}(ch)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
func main() {
jobs := make(chan int, 10)
for i := 1; i <= 6; i++ {
jobs <- i
}
close(jobs)
// fan-out: distribute to 3 workers
workers := fanOut(jobs, 3)
// fan-in: aggregate all worker results
results := fanIn(workers...)
// Collect results
sum := 0
count := 0
for r := range results {
sum += r
count++
}
fmt.Printf("Total %d results, sum = %d\n", count, sum)
}
Resultado:
Worker 2: 1^2 = 1
Worker 0: 3^2 = 9
Worker 1: 2^2 = 4
Worker 1: 5^2 = 25
Worker 0: 4^2 = 16
Worker 2: 6^2 = 36
Total 6 results, sum = 91
(3) fan-out x fan-in
| Modo | Direção | Finalidade |
|---|---|---|
| fan-out | 1 canal → N canais | distribuição de tarefas, computação paralela |
| fan-in | N canais → 1 canal | agregação de resultados, coleta de logs |
flowchart LR
subgraph FanOut [Fan-Out]
J[Jobs Channel] --> W1[Worker 1]
J --> W2[Worker 2]
J --> W3[Worker 3]
end
subgraph FanIn [Fan-In]
W1 --> R[Results Channel]
W2 --> R
W3 --> R
end
style J fill:#e1f5fe
style R fill:#fff3e0
8. Exemplo completo: agregação de cotações de ações de três fontes de dados
// market_aggregator.go
package main
import (
"fmt"
"math/rand"
"sync"
"time"
)
// ---------- Data Types ----------
type Trade struct {
Source string
Symbol string
Price float64
Volume int
Time time.Time
}
type AggregatedQuote struct {
Symbol string
AvgPrice float64
TotalVolume int
Sources int
High float64
Low float64
Time time.Time
}
// ---------- Simulated Exchange Data Source ----------
var symbols = []string{"AAPL", "GOOGL", "MSFT", "AMZN", "TSLA", "META"}
func simulateExchange(name string, out chan<- Trade, done <-chan struct{}) {
for {
select {
case <-done:
return
default:
time.Sleep(time.Duration(100+rand.Intn(300)) * time.Millisecond)
trade := Trade{
Source: name,
Symbol: symbols[rand.Intn(len(symbols))],
Price: 100 + rand.Float64()*200,
Volume: rand.Intn(1000) + 100,
Time: time.Now(),
}
select {
case out <- trade:
case <-done:
return
}
}
}
}
// ---------- Aggregator ----------
type Aggregator struct {
trades chan Trade
done chan struct{}
wg sync.WaitGroup
}
func NewAggregator() *Aggregator {
return &Aggregator{
trades: make(chan Trade, 100),
done: make(chan struct{}),
}
}
func (a *Aggregator) AddSource(name string) {
a.wg.Add(1)
go func() {
defer a.wg.Done()
simulateExchange(name, a.trades, a.done)
}()
fmt.Printf("Added data source: %s\n", name)
}
func (a *Aggregator) Start(interval time.Duration, callback func(map[string]*AggregatedQuote)) {
// Aggregation window
window := make(map[string][]Trade)
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case trade := <-a.trades:
window[trade.Symbol] = append(window[trade.Symbol], trade)
case <-ticker.C:
// Time window reached; compute aggregated quotes
quotes := make(map[string]*AggregatedQuote)
for symbol, trades := range window {
if len(trades) == 0 {
continue
}
var sumPrice float64
var totalVol int
high := trades[0].Price
low := trades[0].Price
for _, t := range trades {
sumPrice += t.Price * float64(t.Volume)
totalVol += t.Volume
if t.Price > high {
high = t.Price
}
if t.Price < low {
low = t.Price
}
}
sources := make(map[string]bool)
for _, t := range trades {
sources[t.Source] = true
}
quotes[symbol] = &AggregatedQuote{
Symbol: symbol,
AvgPrice: sumPrice / float64(totalVol),
TotalVolume: totalVol,
Sources: len(sources),
High: high,
Low: low,
Time: time.Now(),
}
}
callback(quotes)
window = make(map[string][]Trade) // Reset window
case <-a.done:
return
}
}
}
func (a *Aggregator) Stop() {
close(a.done)
a.wg.Wait()
}
// ---------- Main Function ----------
func main() {
aggregator := NewAggregator()
// Add 3 data sources
aggregator.AddSource("NYSE")
aggregator.AddSource("NASDAQ")
aggregator.AddSource("LSE")
fmt.Println("\nStarting aggregation (output every 2 seconds)...\n")
// Start aggregation with 2-second time window
done := make(chan struct{})
go func() {
aggregator.Start(2*time.Second, func(quotes map[string]*AggregatedQuote) {
fmt.Printf("=== Aggregation Report %s ===\n", time.Now().Format("15:04:05"))
for _, q := range quotes {
fmt.Printf("%-6s | $%.2f (avg) | vol=%d | src=%d | H=%.2f L=%.2f\n",
q.Symbol, q.AvgPrice, q.TotalVolume, q.Sources, q.High, q.Low)
}
fmt.Println()
})
close(done)
}()
// Run for 10 seconds then stop
time.Sleep(10 * time.Second)
aggregator.Stop()
<-done
fmt.Println("Quote aggregator stopped")
}
Resultado esperado:
Added data source: NYSE
Added data source: NASDAQ
Added data source: LSE
Starting aggregation (output every 2 seconds)...
=== Aggregation Report 10:00:02 ===
AAPL | $152.34 (avg) | vol=2341 | src=3 | H=165.20 L=142.10
GOOGL | $178.90 (avg) | vol=1567 | src=2 | H=185.00 L=172.30
MSFT | $295.40 (avg) | vol=890 | src=3 | H=301.20 L=288.50
=== Aggregation Report 10:00:04 ===
...
Quote aggregator stopped
select, é necessário se inscrever tanto no case <-done: quanto no case trade := <-a.trades:. Se você se inscrever apenas no canal trades, o sinal Stop não será transmitido. Cada loop for-select deve incluir uma condição de saída.
❓ Perguntas Frequentes
P: O que devo fazer se vários canais estiverem prontos ao mesmo tempo? R: O Go selecionará de forma pseudoaleatória um caso para ser executado. Isso faz parte da especificação da linguagem — para evitar que os desenvolvedores dependam das prioridades dos canais. Se você precisar de prioridades, use duas instruções
selectaninhadas ou uma cláusuladefaultcombinada com verificações manuais.
P: Como o
time.Afterlida com os tempos limite? R: Otime.After(d)retorna um canal que recebe um valor apósdunidades de tempo. Em uma instruçãoselect, usecase <- time.After(d):como o ramo de tempo limite. Observação: Cada chamadaselectparatime.Aftercria um novo temporizador; em um loop, recomenda-se reutilizá-lo comtime.NewTimer.
P: O que
defaultfaz em uma instruçãoselect? R:defaulttorna a instruçãoselectnão-bloqueante — ela executa o blocodefaultimediatamente quando nenhum dos canais está pronto. Adequado para: tentar enviar/receber sem bloqueio, calcular taxas de falha e interrupção/degradação de circuito. Observação: usardefaultem um loopforcria um loop ocupado; normalmente, é necessário adicionartime.Sleepou limitar a frequência.
P: Como faço para sair de um loop for-select usando um canal “done”? R: Crie
done := make(chan struct{}), escutecase <-done:no bloco select e chameclose(done)quando quiser sair. Todas as goroutines que estiverem escutandodonereceberão o sinal de saída simultaneamente.closeé preferível ao envio de um valor — ele pode ser recebido várias vezes.
P: O que são fan-in e fan-out? R: Fan-out = 1 canal de entrada distribuído para N goroutines de trabalho (processamento paralelo); fan-in = N canais agregados em 1 canal de saída (fusão de resultados). A combinação dos dois é um padrão clássico na computação paralela em Go.
P: O
selectpode ser usado para envio? R: Sim. Ocase ch <- value:também é umselect caseválido. Ele fica pronto quando o canal tem espaço (em buffer) ou um receptor (sem buffer). Esse recurso é usado nos modos de limitação de taxa e semáforo.
P: O que acontece com um
select {}vazio? R: Oselect {}ficará bloqueado indefinidamente — pois não há casos nem valor padrão. Isso pode servir como um sinal para “esperar indefinidamente”. No entanto, em código de produção, você deve garantir que haja um caminho de saída; caso contrário, ocorrerá vazamento de goroutines.
P: Uma instrução
breakdentro de um loopfor-selectpode sair do loop? R: Não. Uma instruçãobreakdentro de um blocoselectapenas sai do blocoselect; ela não sai do loopfor. Para sair de um loopfor-select, use um rótulo seguido porbreak label,returnou o canaldone.
📖 Resumo
- o
selectmonitora vários canais simultaneamente e processa aquele que ficar pronto primeiro - Quando vários casos estiverem prontos ao mesmo tempo, escolha um aleatoriamente
time.Afteré usado para controlar o tempo limite, etime.NewTimeré reutilizado em loopsdefaulttorna a operação de seleção não bloqueante, o que é adequado para operações de interrupção de circuito ou de repetição de tentativasdone channelé a maneira padrão de sair de um loop for-select- O “fan-out” distribui tarefas, o “fan-in” agrega resultados, e os dois se combinam para formar um pipeline paralelo
- Todo laço “for-select” deve incluir uma condição de saída
📝 Exercícios
-
Problema básico (Dificuldade ⭐): Escreva um programa que inicie três goroutines, cada uma das quais envie um número inteiro para seu próprio canal. Use
selectpara recebê-los e exibi-los simultaneamente. Adicione um ramo de tempo limite para encerrar o programa após 3 segundos. -
Problema avançado (Dificuldade ⭐⭐): Implemente um cache com tempo limite:
func FetchWithCache(key string, cache map[string]string) string. Primeiro, verifique o canal do cache (tempo limite de 10 ms); se não houver correspondência, consulte o “banco de dados simulado” (500 ms) e preencha o cache com o resultado. Você deve usarselectcombinado comtime.Afterpara implementar tempos de espera em camadas. -
Desafio (Dificuldade: ⭐⭐⭐): Implemente um sistema de agregação de logs: Três fontes de logs (goroutines) geram, cada uma, logs de níveis diferentes (INFO/WARN/ERROR), que são agregados em um único canal por meio de uma operação fan-in e, em seguida, processados por nível usando
select— os logs ERROR acionam um alerta imediato (print), os logs WARN são armazenados em buffer (despejados em lotes de 5) e os logs INFO são gravados em lotes (despejados em lotes de 10). Requisitos: encerramento suave por meio de despejo por tempo limite e do canaldone.