Pool de Workers Concorrente em Go para Processamento de Tarefas HTTP

Para criar um endpoint web e submetê-lo a testes de carga com alto paralelismo, foi implementado um pool de workers no servidor.

A ideia central é separar o recebimento das requisições HTTP do processamento pesado. O handler HTTP apenas enfileira uma tarefa; um conjunto fixo de workers consome a fila de forma concorrente.

Componentes principais:

  • Task: interface com o método Execute() error. Qualquer tipo que implemente esse método pode ser processado pelo pool.
  • Worker: mantém um canal próprio de tarefas, um canal de disponibilidade e um canal de encerramento. O worker registra seu canal de tarefas no pool de workers disponíveis e depois fica bloqueado em select aguardando uma tarefa ou o sinal de parada.
  • Dispatcher: mantém a fila global de tarefas e o pool de canais de workers disponíveis. Sua goroutine principal lê tarefas da fila global, obtém um worker livre e envia a tarefa para o canal desse worker.

Na inicialização do pacote, a fila global é criada e o dispatcher é iniciado. O número de workers é definido por runtime.NumCPU(); em uma máquina com 4 núcleos, por exemplo, são criados 4 workers e, portanto, 4 canais de tarefas concorrentes. A fila pode ter um buffer, como QueueSize = 1024, para absorver picos de requisições.

Implementação do worker:

package workerpool

import "fmt"

type Worker struct {
    ready chan chan Task
    tasks chan Task
    stop  chan struct{}
}

func newWorker(ready chan chan Task) *Worker {
    return &Worker{
        ready: ready,
        tasks: make(chan Task),
        stop:  make(chan struct{}),
    }
}

func (w *Worker) start() {
    go func() {
        for {
            w.ready <- w.tasks

            select {
            case task := <-w.tasks:
                if err := task.Execute(); err != nil {
                    fmt.Println("falha ao executar tarefa:", err)
                }
            case <-w.stop:
                return
            }
        }
    }()
}

func (w *Worker) shutdown() {
    close(w.stop)
}

Implementação do dispatcher e da fila global:

package workerpool

import "runtime"

var (
    WorkerCount = runtime.NumCPU()
    QueueSize   = 1024
    TaskQueue   chan Task
)

type Task interface {
    Execute() error
}

type Dispatcher struct {
    workers int
    ready   chan chan Task
    stop    chan struct{}
    pool    []*Worker
}

func init() {
    runtime.GOMAXPROCS(WorkerCount)
    TaskQueue = make(chan Task, QueueSize)

    dispatcher := NewDispatcher(WorkerCount)
    dispatcher.Start()
}

func NewDispatcher(workers int) *Dispatcher {
    return &Dispatcher{
        workers: workers,
        ready:   make(chan chan Task, workers),
        stop:    make(chan struct{}),
    }
}

func (d *Dispatcher) Start() {
    d.pool = make([]*Worker, 0, d.workers)

    for i := 0; i < d.workers; i++ {
        worker := newWorker(d.ready)
        worker.start()
        d.pool = append(d.pool, worker)
    }

    go d.dispatch()
}

func (d *Dispatcher) dispatch() {
    for {
        select {
        case task := <-TaskQueue:
            workerTasks := <-d.ready
            workerTasks <- task
        case <-d.stop:
            return
        }
    }
}

func (d *Dispatcher) Stop() {
    close(d.stop)

    for _, worker := range d.pool {
        worker.shutdown()
    }
}

Exemplo de uso em um handler HTTP:

package main

import (
    "fmt"
    "net/http"
    "example.com/workerpool"
)

type MobileMessage struct {
    number string
}

func (m *MobileMessage) Execute() error {
    m.number = m.number + "_processado"
    fmt.Println(m.number)
    return nil
}

func handleMobile(w http.ResponseWriter, r *http.Request) {
    defer r.Body.Close()

    if err := r.ParseForm(); err != nil {
        http.Error(w, "formulário inválido", http.StatusBadRequest)
        return
    }

    number := r.PostForm.Get("mobile")
    msg := &MobileMessage{number: number}

    workerpool.TaskQueue <- msg

    w.Header().Set("Content-Type", "application/json")
    w.Write([]byte(`{"status":"ok"}`))
}

func main() {
    http.HandleFunc("/test", handleMobile)

    if err := http.ListenAndServe(":8081", nil); err != nil {
        fmt.Println("falha no servidor:", err)
    }
}

O handler apenas valida o formulário, cria o objeto que implementa workerpool.Task e o envia para workerpool.TaskQueue. O dispatcher entrega a tarefa ao primeiro worker disponível. Como o número de workers acompanha a quantidade de CPUs, o processamento concorrente fica limitado à capacidade de paralelismo da máquina, evitando criar goroutines sem controle para cada requisição.

Tags: go goroutines channels worker-pool net/http

Publicado em 9-18 19:53