Implementando Filas de Mensagens de Alta Concorrência em Go com o Padrão Produtor-Consumidor

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()
}

Tags: go Golang Concurrency producer-consumer Message Queue

Publicado em 7-24 12:52