Documentation
¶
Overview ¶
Package queue e a fila persistente da secao 8 do plano.
A fila mora no Postgres, nao em canal de memoria: "Nunca depender exclusivamente de in-memory channel para jobs criticos". Um processo que morre com itens em canal perde trabalho; um que morre com itens em tabela nao.
O claim usa `FOR UPDATE SKIP LOCKED`, que e o padrao para fila em Postgres: varios dispatchers competem pela mesma tabela sem bloquear uns aos outros e sem entregar o mesmo item duas vezes. A alternativa — SELECT seguido de UPDATE — tem corrida entre as duas instrucoes.
Index ¶
- type Item
- type Queue
- func (q *Queue) Claim(ctx context.Context, worker string, limite int) ([]Item, error)
- func (q *Queue) Done(ctx context.Context, id int64) error
- func (q *Queue) Enqueue(ctx context.Context, runID uuid.UUID, prioridade int, disponivelEm time.Time) error
- func (q *Queue) Recuperar(ctx context.Context, limite time.Duration) ([]Item, error)
- func (q *Queue) Release(ctx context.Context, id int64, atraso time.Duration) error
- func (q *Queue) Tamanho(ctx context.Context) (pendentes, reivindicados int, err error)
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Queue ¶
type Queue struct {
// contains filtered or unexported fields
}
Queue opera sobre queue_items.
func (*Queue) Claim ¶
Claim reivindica ate `limite` itens para este worker.
O limite e como a concorrencia e imposta: o dispatcher pede apenas as vagas que tem livres. Nao existe caminho em que mais itens saiam da fila do que a concorrencia permite, porque quem conta as vagas e quem pede.
func (*Queue) Enqueue ¶
func (q *Queue) Enqueue(ctx context.Context, runID uuid.UUID, prioridade int, disponivelEm time.Time) error
Enqueue coloca um run na fila.
`ON CONFLICT DO NOTHING` na unique de run_id: enfileirar duas vezes o mesmo run e no-op, nao erro. E o comportamento que a secao 29 pede — a operacao tolera repeticao. Enqueue poe o run na fila. `disponivelEm` zero significa AGORA, medido pelo relogio do BANCO.
A diferenca importa: o relogio do processo pode estar alguns milissegundos a frente do relogio do Postgres, e um item gravado com `time.Now()` do aplicativo fica invisivel ate o banco alcanca-lo. Nao e perda — o proximo ciclo pega —, mas e latencia inexplicavel, e foi o que fez um teste de concorrencia entregar 4 itens onde 5 estavam prontos.
func (*Queue) Recuperar ¶
Recuperar devolve a fila os itens reivindicados ha mais tempo que `limite`.
E a rede de seguranca contra worker morto: sem isso, um item reivindicado por um processo que caiu ficaria preso para sempre. Era exatamente o modo de falha das execucoes zumbis que travaram pipelines por 33 dias no sistema anterior.