Padrão Produtor-Consumidor em Go: Fundamentos e Aplicações Concorrentes
O padrão produtor-consumidor é uma pedra angular da programação concorrente, facilitando a comunicação segura entre diferentes partes de um sistema. Na linguagem Go, a implementação desse padrão se torna elegante e eficiente através do uso de canais (channels). Este artigo detalha a construção de vários sistemas de filas concorrentes em Go, desde um modelo básico até uma fila de mensagens de alta performance, explorando compnoentes como produtores, consumidores e distribuição de tarefas.
1. Produtor-Consumidor Básico
A forma mais simples de aplicar o padrão produtor-consumidor envolve um canal para trocar dados. Um goroutine atua como produtor, enviando valores para o canal, enquanto outra goroutine atua como consumidor, lendo esses valores.
package main
import (
"fmt"
"time"
)
// EmitirDados envia uma série de números inteiros para o canal.
func EmitirDados(saida chan<- int) {
for i := 0; i < 10; i++ {
fmt.Printf("Gerando item: %d\n", i)
saida <- i // Envia o item para o canal
time.Sleep(150 * time.Millisecond) // Simula algum trabalho
}
close(saida) // Fecha o canal quando todos os itens forem enviados
}
// ProcessarDados consome itens do canal.
func ProcessarDados(entrada <-chan int, idConsumidor int) {
for item := range entrada {
fmt.Printf("Consumidor %d processou item: %d\n", idConsumidor, item)
time.Sleep(250 * time.Millisecond) // Simula tempo de processamento
}
}
func main() {
// Cria um canal com buffer para itens, com capacidade de 5
canalDeItens := make(chan int, 5)
// Inicia a goroutine produtora
go EmitirDados(canalDeItens)
// Inicia múltiplas goroutines consumidoras
for i := 1; i <= 3; i++ {
go ProcessarDados(canalDeItens, i)
}
// Mantém o programa principal em execução para permitir que goroutines terminem
time.Sleep(6 * time.Second)
fmt.Println("Exemplo básico concluído.")
}
2. Sistema de Fila de Trabalhos
Um sistema de fila de trabalhos (ou task queue) permite que tarefas sejam enfileiradas e processadas por um conjunto de trabalhadores. Isso é ideal para distrbiuir cargas de trabalho e gerenciar tarefas assíncronas.
package main
import (
"fmt"
"time"
)
// Tarefa representa uma unidade de trabalho.
type Tarefa struct {
ID int
Detalhes string
CanalResultado chan string // Canal para reportar o resultado da tarefa
}
// FilaDeTrabalhos gerencia o envio e processamento de tarefas.
type FilaDeTrabalhos struct {
filaInterna chan Tarefa
numeroWorkers int
}
// NovaFilaDeTrabalhos cria uma nova instância da FilaDeTrabalhos.
func NovaFilaDeTrabalhos(workers int, tamanhoBuffer int) *FilaDeTrabalhos {
return &FilaDeTrabalhos{
filaInterna: make(chan Tarefa, tamanhoBuffer),
numeroWorkers: workers,
}
}
// IniciarWorkers lança as goroutines dos trabalhadores.
func (ft *FilaDeTrabalhos) IniciarWorkers() {
for i := 0; i < ft.numeroWorkers; i++ {
go ft.executarWorker(i + 1)
}
}
// executarWorker é a lógica de cada goroutine de trabalhador.
func (ft *FilaDeTrabalhos) executarWorker(idWorker int) {
for tarefa := range ft.filaInterna {
fmt.Printf("Worker %d está processando tarefa %d: %s\n", idWorker, tarefa.ID, tarefa.Detalhes)
time.Sleep(700 * time.Millisecond) // Simula processamento intensivo
tarefa.CanalResultado <- fmt.Sprintf("Tarefa %d (Worker %d) finalizada com sucesso.", tarefa.ID, idWorker)
}
}
// EnviarTrabalho adiciona uma tarefa à fila.
func (ft *FilaDeTrabalhos) EnviarTrabalho(tarefa Tarefa) {
ft.filaInterna <- tarefa
}
// FecharFila encerra o canal de tarefas, sinalizando para os workers pararem.
func (ft *FilaDeTrabalhos) FecharFila() {
close(ft.filaInterna)
}
func main() {
// Configura uma fila com 3 workers e capacidade de 10 tarefas
gerenciador := NovaFilaDeTrabalhos(3, 10)
gerenciador.IniciarWorkers()
// Envio de tarefas e coleta assíncrona de resultados
for i := 1; i <= 10; i++ {
canalDeRetorno := make(chan string, 1)
tarefaAtual := Tarefa{
ID: i,
Detalhes: fmt.Sprintf("Processar dados do lote %d", i),
CanalResultado: canalDeRetorno,
}
gerenciador.EnviarTrabalho(tarefaAtual)
// Goroutine para coletar o resultado da tarefa específica
go func(idTarefa int, ch <-chan string) {
resultado := <-ch
fmt.Printf("Resultado da tarefa %d: %s\n", idTarefa, resultado)
}(i, canalDeRetorno)
}
time.Sleep(12 * time.Second) // Aguarda o processamento das tarefas
gerenciador.FecharFila()
fmt.Println("Sistema de fila de trabalhos concluído.")
}
3. Fila de Prioridades
Para cenários onde algumas tarefas são mais críticas que outras, uma fila de prioridades é essencial. Em Go, isso pode ser implementado usando o pacote container/heap para manter as tarefas ordenadas por prioridade.
package main
import (
"container/heap"
"context" // Importar context para usar o contexto.Context e context.CancelFunc
"fmt"
"sync"
"time"
)
// ItemPrioritario representa uma tarefa com um nível de prioridade.
type ItemPrioritario struct {
Identificador int
Descricao string
Prioridade int // Maior valor = maior prioridade
CanalFeedback chan string
}
// HeapDePrioridades implementa a interface heap.Interface para ItemPrioritario.
type HeapDePrioridades []*ItemPrioritario
func (hdp HeapDePrioridades) Len() int { return len(hdp) }
// Less retorna verdadeiro se o item 'i' tiver maior prioridade que o item 'j'.
// Estamos usando uma min-heap por padrão, então para maior prioridade na frente, invertemos a lógica.
func (hdp HeapDePrioridades) Less(i, j int) bool {
return hdp[i].Prioridade > hdp[j].Prioridade // Os com maior prioridade vêm primeiro
}
func (hdp HeapDePrioridades) Swap(i, j int) {
hdp[i], hdp[j] = hdp[j], hdp[i]
}
func (hdp *HeapDePrioridades) Push(x interface{}) {
item := x.(*ItemPrioritario)
*hdp = append(*hdp, item)
}
func (hdp *HeapDePrioridades) Pop() interface{} {
old := *hdp
n := len(old)
item := old[n-1]
*hdp = old[0 : n-1]
return item
}
// GerenciadorDePrioridades coordena a fila de tarefas prioritárias e os workers.
type GerenciadorDePrioridades struct {
filaDeEntrada chan *ItemPrioritario
numProcessadores int
filaInterna HeapDePrioridades
mutex sync.Mutex // Protege o acesso à filaInterna
ctx context.Context
cancel context.CancelFunc
}
// NovoGerenciadorDePrioridades cria um novo gerenciador de fila de prioridades.
func NovoGerenciadorDePrioridades(processadores int, capacidadeEntrada int) *GerenciadorDePrioridades {
ctx, cancel := context.WithCancel(context.Background())
gp := &GerenciadorDePrioridades{
filaDeEntrada: make(chan *ItemPrioritario, capacidadeEntrada),
numProcessadores: processadores,
filaInterna: make(HeapDePrioridades, 0),
ctx: ctx,
cancel: cancel,
}
heap.Init(&gp.filaInterna) // Inicializa a estrutura de heap
return gp
}
// IniciarOperacao lança o agendador e os workers.
func (gp *GerenciadorDePrioridades) IniciarOperacao() {
go gp.rotinaAgendadora()
for i := 0; i < gp.numProcessadores; i++ {
go gp.rotinaDeProcessamento(i + 1)
}
}
// rotinaAgendadora recebe tarefas e as adiciona ao heap de prioridades.
func (gp *GerenciadorDePrioridades) rotinaAgendadora() {
for {
select {
case item := <-gp.filaDeEntrada:
gp.mutex.Lock()
heap.Push(&gp.filaInterna, item)
gp.mutex.Unlock()
case <-gp.ctx.Done():
fmt.Println("Agendador encerrado.")
return
}
}
}
// rotinaDeProcessamento extrai e executa tarefas com base na prioridade.
func (gp *GerenciadorDePrioridades) rotinaDeProcessamento(idWorker int) {
for {
select {
case <-gp.ctx.Done():
fmt.Printf("Worker %d encerrado.\n", idWorker)
return
default:
gp.mutex.Lock()
if gp.filaInterna.Len() > 0 {
tarefa := heap.Pop(&gp.filaInterna).(*ItemPrioritario)
gp.mutex.Unlock()
fmt.Printf("Worker %d processando tarefa %d (Prioridade: %d)\n",
idWorker, tarefa.Identificador, tarefa.Prioridade)
time.Sleep(600 * time.Millisecond) // Simula o trabalho
tarefa.CanalFeedback <- fmt.Sprintf("Tarefa %d finalizada.", tarefa.Identificador)
} else {
gp.mutex.Unlock()
time.Sleep(50 * time.Millisecond) // Espera um pouco se a fila estiver vazia
}
}
}
}
// EnviarItem insere um novo item na fila de entrada.
func (gp *GerenciadorDePrioridades) EnviarItem(item *ItemPrioritario) {
gp.filaDeEntrada <- item
}
// Encerrar gracefully shuts down the queue system.
func (gp *GerenciadorDePrioridades) Encerrar() {
close(gp.filaDeEntrada) // Fecha o canal de entrada
gp.cancel() // Sinaliza para agendador e workers pararem
time.Sleep(1 * time.Second) // Give workers some time to finish
}
func main() {
queue := NovoGerenciadorDePrioridades(3, 20)
queue.IniciarOperacao()
// Submete tarefas com prioridades variadas
niveisPrioridade := []int{3, 1, 5, 2, 4, 10, 6, 7, 8, 9}
for i, p := range niveisPrioridade {
canalDeFeedback := make(chan string, 1)
item := &ItemPrioritario{
Identificador: i + 1,
Descricao: fmt.Sprintf("Pacote de dados %d", i+1),
Prioridade: p,
CanalFeedback: canalDeFeedback,
}
queue.EnviarItem(item)
go func(id int, ch <-chan string) {
msg := <-ch
fmt.Printf("Feedback da tarefa %d: %s\n", id, msg)
}(item.Identificador, canalDeFeedback)
}
time.Sleep(15 * time.Second) // Aguarda o processamento das tarefas
queue.Encerrar()
fmt.Println("Fila de prioridades concluída.")
}
4. Sistema de Fila de Mensagens
Um sistema de fila de mensagens (ou message broker simplificado) permite a publicação de mensagens em tópicos e a subscrição por múltiplos consumidores. Isso é fundamental para arquiteturas de microsserviços e desacoplamento de componentes.
package main
import (
"context"
"fmt"
"sync"
"time"
)
// Evento representa uma mensagem a ser transmitida.
type Evento struct {
ID int
Topico string
Conteudo []byte
Timestamp time.Time
}
// BrokerDeEventos gerencia tópicos e assinaturas.
type BrokerDeEventos struct {
canaisPorTopico map[string]chan Evento
assinantes map[string][]chan Evento
mu sync.RWMutex // Protege os mapas de tópicos e assinantes
contexto context.Context
cancelFunc context.CancelFunc
}
// NovoBrokerDeEventos cria uma nova instância do BrokerDeEventos.
func NovoBrokerDeEventos() *BrokerDeEventos {
ctx, cancel := context.WithCancel(context.Background())
return &BrokerDeEventos{
canaisPorTopico: make(map[string]chan Evento),
assinantes: make(map[string][]chan Evento),
contexto: ctx,
cancelFunc: cancel,
}
}
// CriarCanalDeTopico inicializa um novo tópico e seu distribuidor.
func (be *BrokerDeEventos) CriarCanalDeTopico(nomeTopico string, capacidade int) {
be.mu.Lock()
defer be.mu.Unlock()
if _, existe := be.canaisPorTopico[nomeTopico]; !existe {
be.canaisPorTopico[nomeTopico] = make(chan Evento, capacidade)
be.assinantes[nomeTopico] = make([]chan Evento, 0)
go be.rotinaDeDistribuicao(nomeTopico) // Inicia o distribuidor para este tópico
}
}
// rotinaDeDistribuicao lê eventos de um tópico e os envia para todos os assinantes.
func (be *BrokerDeEventos) rotinaDeDistribuicao(nomeTopico string) {
canalDoTopico := be.canaisPorTopico[nomeTopico]
for {
select {
case evento := <-canalDoTopico:
be.mu.RLock()
consumidoresDoTopico := be.assinantes[nomeTopico]
be.mu.RUnlock()
// Transmite o evento para todos os canais de consumidores inscritos
for _, canalConsumidor := range consumidoresDoTopico {
select {
case canalConsumidor <- evento:
// Evento enviado com sucesso
default:
// Canal do consumidor cheio, ignora para não bloquear o distribuidor
fmt.Printf("Alerta: Canal de consumidor cheio para tópico '%s', evento %d descartado.\n", nomeTopico, evento.ID)
}
}
case <-be.contexto.Done():
fmt.Printf("Distribuidor do tópico '%s' encerrado.\n", nomeTopico)
return
}
}
}
// PublicarEvento envia uma mensagem para um tópico específico.
func (be *BrokerDeEventos) PublicarEvento(nomeTopico string, evento Evento) error {
be.mu.RLock()
canalDoTopico, existe := be.canaisPorTopico[nomeTopico]
be.mu.RUnlock()
if !existe {
return fmt.Errorf("tópico '%s' não existe", nomeTopico)
}
select {
case canalDoTopico <- evento:
return nil
case <-be.contexto.Done():
return be.contexto.Err()
case <-time.After(100 * time.Millisecond): // Timeout para evitar bloqueio eterno se o canal estiver cheio
return fmt.Errorf("falha ao publicar no tópico '%s': canal cheio ou timeout", nomeTopico)
}
}
// AssinarTopico permite que um consumidor receba mensagens de um tópico.
func (be *BrokerDeEventos) AssinarTopico(nomeTopico string) <-chan Evento {
be.mu.Lock()
defer be.mu.Unlock()
canalDeConsumo := make(chan Evento, 20) // Canal do consumidor com buffer
be.assinantes[nomeTopico] = append(be.assinantes[nomeTopico], canalDeConsumo)
return canalDeConsumo
}
// EncerrarBroker sinaliza para todos os componentes pararem.
func (be *BrokerDeEventos) EncerrarBroker() {
be.cancelFunc()
time.Sleep(500 * time.Millisecond) // Pequena pausa para garantir que os goroutines recebam o sinal
}
func main() {
broker := NovoBrokerDeEventos()
broker.CriarCanalDeTopico("noticias", 100)
broker.CriarCanalDeTopico("clima", 50)
// Assinantes dos tópicos
canalNoticias := broker.AssinarTopico("noticias")
canalClima := broker.AssinarTopico("clima")
// Goroutine para consumir notícias
go func() {
for evt := range canalNoticias {
fmt.Printf("[NOTICIAS] Recebido: %s (ID: %d)\n", string(evt.Conteudo), evt.ID)
}
fmt.Println("Consumidor de notícias finalizado.")
}()
// Goroutine para consumir clima
go func() {
for evt := range canalClima {
fmt.Printf("[CLIMA] Recebido: %s (ID: %d)\n", string(evt.Conteudo), evt.ID)
}
fmt.Println("Consumidor de clima finalizado.")
}()
// Produtor de notícias
go func() {
for i := 0; i < 8; i++ {
evento := Evento{
ID: i + 1,
Topico: "noticias",
Conteudo: []byte(fmt.Sprintf("Atualização de notícia #%d", i+1)),
Timestamp: time.Now(),
}
if err := broker.PublicarEvento("noticias", evento); err != nil {
fmt.Printf("Erro ao publicar notícia: %v\n", err)
}
time.Sleep(400 * time.Millisecond)
}
fmt.Println("Produtor de notícias finalizado.")
}()
// Produtor de clima
go func() {
for i := 0; i < 6; i++ {
evento := Evento{
ID: i + 101,
Topico: "clima",
Conteudo: []byte(fmt.Sprintf("Previsão do tempo para hoje #%d", i+1)),
Timestamp: time.Now(),
}
if err := broker.PublicarEvento("clima", evento); err != nil {
fmt.Printf("Erro ao publicar clima: %v\n", err)
}
time.Sleep(600 * time.Millisecond)
}
fmt.Println("Produtor de clima finalizado.")
}()
time.Sleep(10 * time.Second) // Permite que os produtores e consumidores trabalhem
broker.EncerrarBroker()
fmt.Println("Sistema de fila de mensagens encerrado.")
}
5. Padrão de Piscina de Trabalhadores (Worker Pool)
O padrão worker pool é uma forma eficeinte de gerenciar e limitar a concorrência. Um número fixo de goroutines (trabalhadores) aguarda por tarefas em um canal, processando-as à medida que chegam. Isso evita a sobrecarga do sistema com a criação excessiva de goroutines.
package main
import (
"fmt"
"sync"
"time"
)
// Trabalho representa uma tarefa a ser executada.
type Trabalho struct {
Identificador int
ConteudoDados string
}
// PiscinaDeTrabalhadores gerencia um grupo de workers.
type PiscinaDeTrabalhadores struct {
canalTrabalhos chan Trabalho
canalResultados chan string
numWorkers int
wg sync.WaitGroup // Para aguardar a finalização de todos os workers
mu sync.Mutex // Protege o contador de workers ativos
workersAtivos int
}
// NovaPiscinaDeTrabalhadores cria uma nova piscina.
func NovaPiscinaDeTrabalhadores(totalWorkers int, tamanhoFilaTrabalhos int) *PiscinaDeTrabalhadores {
return &PiscinaDeTrabalhadores{
canalTrabalhos: make(chan Trabalho, tamanhoFilaTrabalhos),
canalResultados: make(chan string, tamanhoFilaTrabalhos),
numWorkers: totalWorkers,
}
}
// IniciarPiscina lança as goroutines dos trabalhadores.
func (pdt *PiscinaDeTrabalhadores) IniciarPiscina() {
for i := 0; i < pdt.numWorkers; i++ {
pdt.wg.Add(1)
go pdt.rotinaDeWorker(i + 1)
}
}
// rotinaDeWorker é a lógica que cada trabalhador executa.
func (pdt *PiscinaDeTrabalhadores) rotinaDeWorker(idWorker int) {
defer pdt.wg.Done()
pdt.mu.Lock()
pdt.workersAtivos++
pdt.mu.Unlock()
for trabalho := range pdt.canalTrabalhos {
fmt.Printf("Worker %d iniciando trabalho %d: %s\n", idWorker, trabalho.Identificador, trabalho.ConteudoDados)
time.Sleep(750 * time.Millisecond) // Simula o processamento
pdt.canalResultados <- fmt.Sprintf("Trabalho %d concluído pelo Worker %d.", trabalho.Identificador, idWorker)
}
pdt.mu.Lock()
pdt.workersAtivos--
pdt.mu.Unlock()
fmt.Printf("Worker %d encerrado.\n", idWorker)
}
// EnviarTrabalho envia um trabalho para a piscina.
func (pdt *PiscinaDeTrabalhadores) EnviarTrabalho(trabalho Trabalho) {
pdt.canalTrabalhos <- trabalho
}
// EncerrarPiscina fecha o canal de trabalhos e espera por todos os workers.
func (pdt *PiscinaDeTrabalhadores) EncerrarPiscina() {
close(pdt.canalTrabalhos) // Sinaliza para os workers pararem de receber novos trabalhos
pdt.wg.Wait() // Aguarda a conclusão de todos os workers
close(pdt.canalResultados) // Fecha o canal de resultados após todos os trabalhos serem processados
fmt.Println("Piscina de trabalhadores encerrada.")
}
// ObterNumeroWorkersAtivos retorna o número de workers atualmente ativos.
func (pdt *PiscinaDeTrabalhadores) ObterNumeroWorkersAtivos() int {
pdt.mu.Lock()
defer pdt.mu.Unlock()
return pdt.workersAtivos
}
func main() {
// Cria uma piscina com 4 workers e uma fila de trabalhos de 15 posições
piscina := NovaPiscinaDeTrabalhadores(4, 15)
piscina.IniciarPiscina()
// Envio de 18 trabalhos para a piscina
for i := 1; i <= 18; i++ {
piscina.EnviarTrabalho(Trabalho{
Identificador: i,
ConteudoDados: fmt.Sprintf("Item de processamento #%d", i),
})
}
// Goroutine para coletar e exibir os resultados
go func() {
for resultado := range piscina.canalResultados {
fmt.Println("Resultado coletado:", resultado)
}
fmt.Println("Coletor de resultados finalizado.")
}()
time.Sleep(15 * time.Second) // Aguarda um tempo para que os trabalhos sejam processados
piscina.EncerrarPiscina()
}