Intensivão Golang: concorrência, resiliência e sistemas distribuídos em 30 minutos

Um roteiro de 30 minutos para reativar Go aplicado a serviços de produção, ingestão de telemetria e sistemas distribuídos.

Este roteiro serve para quem já trabalha com backend e quer reativar Go para uma conversa técnica ou um serviço de produção. O foco está nas decisões que mantêm um sistema previsível sob carga: concorrência limitada, cancelamento, filas finitas, idempotência e observabilidade.

O recorte usa Go 1.26. Não é uma introdução à linguagem. Passe rápido pelos fundamentos e retenha os pontos que mudam o desenho de um consumer, uma API ou um pipeline de telemetria.

Roteiro de 30 minutos

TempoBlocoPrioridade
0-4 minTipos, structs, interfaces e errosRevisão rápida
4-10 minGoroutines, channels, select e contextoAlta
10-17 minLimites de concorrência e backpressureMáxima
17-23 minPipeline IoT resilienteMáxima
23-26 minRuntime, memória e profilingAlta
26-30 minArquitetura e perguntas de entrevistaMáxima

Fundamentos que aparecem em produção

Go favorece composição, contratos pequenos e fluxo explícito. Não há herança de classes nem exceções como mecanismo normal de controle. Um tipo simples pode carregar sua própria validação:

package telemetry

import (
	"errors"
	"fmt"
	"time"
)

var ErrOutOfRange = errors.New("reading out of range")

type Reading struct {
	DeviceID   string    `json:"device_id"`
	Sequence   uint64    `json:"sequence"`
	ObservedAt time.Time `json:"observed_at"`
	Value      float64   `json:"value"`
}

func (r Reading) Validate() error {
	if r.DeviceID == "" {
		return errors.New("device_id is required")
	}
	if r.Value < -100 || r.Value > 250 {
		return fmt.Errorf("%w: %.2f", ErrOutOfRange, r.Value)
	}
	return nil
}

Alguns detalhes evitam erros silenciosos:

  • O zero value costuma ser utilizável. Prefira tipos cujo estado inicial seja válido quando isso não esconder uma regra de negócio.
  • Slice é uma visão sobre um array. Cópias podem compartilhar o mesmo backing array; append pode reutilizá-lo ou alocar outro.
  • map não tem ordem de iteração e não suporta leitura e escrita concorrentes sem sincronização.
  • string contém bytes imutáveis, normalmente UTF-8. len conta bytes; range decodifica runes.
  • defer executa em LIFO, mas avalia os argumentos quando é registrado.

Interfaces são satisfeitas implicitamente. Defina a interface pequena no pacote que a consome, em vez de exportar um contrato grande ao lado da implementação. Erros são valores: acrescente contexto com %w e inspecione a causa com errors.Is ou errors.As.

if err := reading.Validate(); err != nil {
	if errors.Is(err, ErrOutOfRange) {
		return sendToDLQ(reading, err)
	}
	return fmt.Errorf("validate reading: %w", err)
}

panic fica para invariantes quebradas ou falha irrecuperável de inicialização. recover só alcança um panic na mesma goroutine e pertence a fronteiras controladas, como middleware.

Goroutines, channels e cancelamento

Uma goroutine não é uma thread dedicada. O runtime a agenda sobre threads do sistema operacional. Iniciar uma goroutine sem saber quem a cancela ou espera cria risco de leak.

Channels transportam trabalho ou propriedade. Um mutex protege estado compartilhado. Um buffer apenas absorve uma diferença temporária de velocidade; não cria capacidade infinita.

jobs := make(chan Reading, 128)

go func() {
	defer close(jobs) // quem produz fecha
	for _, reading := range batch {
		jobs <- reading
	}
}()

for reading := range jobs {
	if err := process(reading); err != nil {
		// tratar ou registrar o erro da mensagem
	}
}

O produtor fecha o channel quando não haverá novo envio. Enviar em channel fechado ou fechá-lo duas vezes causa panic. Receber de um channel fechado retorna o zero value e ok == false. Um channel nil bloqueia para sempre; dentro de select, ele desabilita o caso.

select combina envio, cancelamento e política de saturação. Sem default, a operação espera espaço ou cancelamento. Com default, ela rejeita imediatamente quando a fila está cheia:

func enqueue(ctx context.Context, jobs chan<- Reading, reading Reading) error {
	select {
	case jobs <- reading:
		return nil
	case <-ctx.Done():
		return context.Cause(ctx)
	}
}

context.Context carrega cancelamento, deadline e metadados estritamente ligados à requisição. Receba-o como primeiro argumento, propague-o, chame todo cancel retornado e não guarde contexto em struct. Cancelar não encerra uma goroutine à força: os loops e operações bloqueantes precisam observar ctx.Done().

Concorrência limitada antes da pressão de memória

Uma goroutine por mensagem parece barata até um downstream ficar lento. A fila cresce, o heap cresce, o GC trabalha mais e o processo pode cair antes de a CPU parecer saturada.

Para tarefas independentes que falham juntas, errgroup oferece espera, propagação do primeiro erro e cancelamento compartilhado. SetLimit estabelece o teto de concorrência.

func ProcessBatch(ctx context.Context, batch []Reading) error {
	g, ctx := errgroup.WithContext(ctx)
	g.SetLimit(16)

	for _, reading := range batch {
		reading := reading
		g.Go(func() error {
			if err := processOne(ctx, reading); err != nil {
				return fmt.Errorf("device %s: %w", reading.DeviceID, err)
			}
			return nil
		})
	}

	return g.Wait()
}

Classifique o erro antes de decidir o que fazer. Um banco indisponível pode justificar cancelar o lote. Um payload inválido, duplicado ou fora do schema deve ir para quarentena ou DLQ, sem derrubar o consumer inteiro.

Escolha a primitiva pela propriedade que precisa preservar:

NecessidadePrimitiva
Contador ou flag independenteatomic tipado
Invariante entre vários campossync.Mutex
Leitura frequente e escrita curtasync.RWMutex, depois de medir
Transferir trabalho ou ownershipchannel
Inicialização únicasync.Once
Esperar tarefas sem errosync.WaitGroup
Esperar tarefas com erro e cancelamentoerrgroup

Não copie mutex depois do primeiro uso e não mantenha lock durante I/O remoto. RWMutex não é uma melhoria automática para uma seção crítica pequena.

Backpressure é uma decisão de produto e operação. Se a ingestão recebe 50 mil mensagens por segundo e a persistência confirma 20 mil, acumular o restante em memória apenas muda o incidente de lugar.

PolíticaConsequência
Bloquear produtorAumenta latência e preserva dados quando o protocolo aceita desacelerar
Rejeitar com erroExige retry e idempotência no cliente
Pausar consumo ou ACKMantém backlog no broker durável
Descartar dados antigosPreserva frescor quando histórico não importa
Agregar ou downsampleReduz resolução para aliviar a carga
Persistir em discoEvita perda, com custo operacional adicional

Defina tamanho de buffer, métrica de ocupação, timeout e ação de saturação. Buffer sem política não é estratégia de capacidade.

Pipeline IoT que tolera reentrega

Uma separação comum é:

dispositivo -> MQTT/broker -> ingestão Go -> stream -> processadores -> armazenamento
                                  \-> DLQ        \-> estado atual

MQTT atende bem à borda e às conexões dos dispositivos. Um stream como Kafka atende retenção, replay e particionamento interno. gRPC é RPC interno tipado; WebSocket atende atualização de dashboards. Nenhum deles substitui os outros automaticamente.

Fan-out de workers quebra ordem global. Em telemetria, o requisito costuma ser ordem por dispositivo. Particione por uma chave estável, como hash(device_id) % N, e processe cada partição de forma sequencial. Guarde observed_at, ingested_at, sequence, event_id e, quando existir, boot_id. O relógio do dispositivo pode estar errado ou reiniciar.

Projete a cadeia para entrega at least once. Um fluxo seguro recebe o evento, valida envelope e versão, verifica a chave de idempotência, grava efeito e marcador de deduplicação na mesma transação quando possível e só então confirma a mensagem. Para publicação após uma atualização de banco, uma transactional outbox elimina a janela entre confirmar a transação e publicar o evento. Duplicatas ainda podem ocorrer, portanto o consumidor continua idempotente.

Retry serve para falha transitória. Payload inválido e regra de negócio rejeitada não melhoram com novas tentativas. Use limite, budget total e jitter para evitar que todos os pods repitam juntos:

func retry(ctx context.Context, max int, fn func(context.Context) error) error {
	var err error
	for attempt := 0; attempt < max; attempt++ {
		if err = fn(ctx); err == nil {
			return nil
		}

		base := min(100*time.Millisecond<<attempt, 5*time.Second)
		wait := time.Duration(rand.Int64N(int64(base) + 1))
		timer := time.NewTimer(wait)
		select {
		case <-timer.C:
		case <-ctx.Done():
			if !timer.Stop() {
				select {
				case <-timer.C:
				default:
				}
			}
			return context.Cause(ctx)
		}
	}
	return fmt.Errorf("retry exhausted: %w", err)
}

Na borda, imponha TLS, identidade individual por dispositivo, autorização por tópico, limites de payload e validação antes de alocar estruturas grandes. Rotação, revogação, sequência, nonce e janela temporal entram quando o protocolo de negócio precisa resistir a replay. Não registre credenciais ou payloads sensíveis sem critério.

Kubernetes, shutdown e observabilidade

Em SIGTERM, um consumer deve sair de readiness, parar de buscar trabalho, drenar o que já está em voo dentro do grace period e confirmar apenas o que terminou. O restante volta ao broker. Feche produtores, conexões e telemetria por último.

Liveness responde se o processo progride; não a faça depender de cada serviço externo. Readiness responde se o pod pode receber trabalho agora e pode falhar por overload ou perda de uma dependência obrigatória. Startup protege inicializações lentas.

Para escalar consumers, CPU isolada costuma ser um sinal fraco. Observe lag, idade da mensagem mais antiga, taxa de entrada, duração do processamento e ocupação do pool.

Logs estruturados devem carregar a correlação necessária para investigar uma leitura sem despejar o payload inteiro:

logger.InfoContext(ctx, "reading persisted",
	"device_id", reading.DeviceID,
	"sequence", reading.Sequence,
	"latency_ms", elapsed.Milliseconds(),
)

Métricas úteis incluem throughput, erros por classe, p50/p95/p99, lag, idade do evento, espera e ocupação dos workers, retries, DLQ, duplicatas, goroutines, heap e pausas de GC. Use tracing amostrado para atravessar ingestão, stream e persistência. Traçar cada leitura de alta frequência pode custar mais do que a investigação que ele pretende facilitar.

Runtime e performance

Uma data race ocorre quando acessos concorrentes à mesma posição de memória incluem escrita e não são ordenados por sincronização. Envio em channel, Mutex.Unlock seguido da aquisição do mesmo lock e operações atômicas fornecem relações de ordem relevantes. O race detector ajuda, mas só encontra caminhos executados:

go test -race ./...

O scheduler trabalha com G (goroutine), M (thread do sistema) e P (recurso lógico de execução). GOMAXPROCS limita quantos Ps executam código Go simultaneamente, não quantas goroutines podem existir. Goroutines são leves, mas stacks, referências e scheduling têm custo. Crie trabalho com limite.

Antes de otimizar, estabeleça um SLO, reproduza a carga e meça. pprof mostra CPU, heap, bloqueio e contenção; go tool trace ajuda a observar scheduler e latências de bloqueio. Benchmarks isolam a mudança.

go test ./...
go test -race ./...
go test -bench=. -benchmem ./...
go test -run=^$ -bench=BenchmarkDecode -cpuprofile=cpu.out -memprofile=mem.out ./...
go tool pprof cpu.out
go build -gcflags="-m=2" ./...

Pré-alocar um slice com capacidade conhecida reduz realocações. Evite conversões repetidas entre string e []byte em caminho quente. sync.Pool é um cache oportunista para temporários, não uma reserva de capacidade. Pooling sem profile aumenta a complexidade e pode reter memória sem ganho.

Generics funcionam bem para algoritmos que preservam as mesmas operações entre tipos. Interfaces continuam melhores para comportamento dinâmico e fronteiras. Não transforme DTOs e serviços de domínio em abstrações genéricas sem um ganho claro.

Perguntas que valem uma entrevista

Goroutine é thread? Não. É uma unidade leve gerenciada pelo runtime e multiplexada sobre threads do sistema operacional.

Channel ou mutex? Use channel para transferir trabalho ou ownership; use mutex para proteger estado compartilhado e invariantes.

Quem fecha um channel? O produtor que sabe que não haverá mais envios. O receiver normalmente não fecha.

Como preservar ordem com vários workers? Não preserve ordem global sem necessidade. Particione por uma chave, como device_id, e mantenha cada partição sequencial.

Como lidar com duplicatas? Use uma chave de idempotência estável, deduplicação transacional quando possível e operações naturalmente idempotentes, como upsert com versão ou sequência.

Como parar um consumer no Kubernetes? Receba SIGTERM, remova readiness, pare fetch, drene o trabalho em voo dentro do grace period, confirme apenas o concluído e feche recursos.

Como investigar latência? Separe fila, processamento e dependências. Compare p95/p99, lag e saturação; depois use tracing, pprof e go tool trace para testar uma hipótese.

Checklist de produção

  • Quem é dono de cada goroutine e como ela termina?
  • Onde o contexto propaga cancelamento e deadline?
  • Qual é o limite de concorrência e o comportamento na saturação?
  • A ordem necessária é por chave ou realmente global?
  • Quando ocorre ACK e como a mutação resiste a reentrega?
  • Quais erros vão para retry, DLQ ou quarentena?
  • O que diferencia horário do dispositivo de horário de ingestão?
  • Como TLS, identidade e autorização impedem publicação indevida?
  • O shutdown para de aceitar trabalho antes de encerrar o processo?
  • Quais métricas mostram lag, p99, heap, GC e contenção?

Referências oficiais