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étodoExecute() 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 emselectaguardando 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.