Áttekintés
A queue egy önálló Go modul (github.com/gp-system/queue): a Valkey/asynq alapréteg, amire az async kategória minden további tagja épül: az események, az outbox, az ütemezés és a kit worker chassisa. Egy helyen fogja össze a connection-konfigot és a task-beküldést; a szerver oldala az asynq, ami a worker-feldolgozást és a monitorozást adja. A felsőbb szintű opciókat, az event-borítékot és a task-típus-elnevezést a következő oldal, a Taskok tárgyalja; ez az oldal a kapcsolódás és a beküldés szintjén marad.
import "github.com/gp-system/queue"
A mindennapokban ritkán hívod közvetlenül: az eventeket a dispatcher, a jobokat az ütemező viszi a queue-ba. De az itteni fogalmak (config, kliens, a boríték a következő oldalon) köszönnek vissza minden felsőbb rétegben, ezért érdemes itt kezdeni.
queuecentralizál, nem absztrahál: az asynq típusai (*asynq.Task, *asynq.TaskInfo, asynq.Option) látszanak a publikus API-ban. Ha az event/scheduler réteg nem illik a feladatra, az alkalmazáskód ugyanezen a kliensen át nyugodtan használhatja az asynq-ot közvetlenül. Mivel semmi nincs elrejtve, az asynq ökoszisztéma eszközei (pl. az Asynqmon webes UI, a Horizon dashboard megfelelője) ugyanarra a Valkeyre kötve azonnal működnek.Telepítés
go get github.com/gp-system/queue@v0.1.0 # a kit is ezt a taget használja
go get github.com/gp-system/queue@latest
Go 1.25+. A queue az errs-től, a hibiken/asynq-tól és a redis/go-redis/v9-től függ, plusz az OTel tracing API-jától a producer span miatt.
Konfiguráció
A queue.Config az egyetlen Valkey-kapcsolatleírás: a kliens, a worker, az ütemező és az outbox relay is ezt kapja. Prefix alá komponálva használd:
type Config struct {
Valkey queue.Config `envPrefix:"VALKEY_"`
}
| Változó | Default | Jelentés |
|---|---|---|
VALKEY_ADDR | localhost:6379 | Valkey host:port |
VALKEY_PASSWORD | (üres) | jelszó; üresen nincs auth |
VALKEY_DB | 0 | logikai Valkey-adatbázis |
A configból két irányba nyílik kapcsolat, mindkét library ugyanabból az egy forrásból konfigurálódik:
cfg.ValkeyConnOpt()→asynq.RedisClientOpt(asynq kliens, szerver, scheduler),cfg.ValkeyOptions()→*redis.Options(közvetlen go-redis hívások, pl. az ütemező lease-e).
A worker Config-ja már tartalmazza VALKEY_ prefixszel. A generált projektben külön bedrótozni nem kell.
Client
client, err := queue.NewClient(ctx, cfg.Valkey) // vagy queue.MustNewClient main()-ben
defer client.Close()
info, err := client.Enqueue(ctx, task, queue.OnQueue("mail"))
A NewClient megnyitja a Valkey-kapcsolatot, és konstruáláskor PING-gel ellenőrzi: ugyanaz a konvenció, mint a pg.NewPool-nál, a rossz cím induláskor bukik, nem az első enqueue-nál. A MustNewClient a panikoló változat main()-be.
Az Enqueue(ctx, *asynq.Task, opts...) beadja a taskot a Valkeybe, és *asynq.TaskInfo-t ad vissza. Két dolgot tesz hozzá a nyers asynq-hoz:
- Producer span: minden enqueue OpenTelemetry spant nyit (
enqueue <task-type>), így a queue-ba adás látszik a trace-ben (lásd trace-propagáció a következő oldalon). queue.ErrDuplicate: ha egyTaskIDvagyUniquemegkötés miatt a task már bent van, az asynq kétféle hibája helyett egyetlen szentinel hibát kapsz. Az idempotens kézbesítésre építő hívók (outbox relay, event fan-out) ezt sikernek tekintik:
if _, err := client.Enqueue(ctx, task, queue.TaskID(id)); err != nil {
if errors.Is(err, queue.ErrDuplicate) {
return nil // már bent van, pont ezt akartuk
}
return err
}
A Close() a megosztott Valkey-kapcsolatot zárja (az asynq kliens ebből épül, és nem birtokolja). A worker ezt a shutdown-sorrend részeként magától megteszi. Az Enqueue további opcióit (OnQueue, MaxRetry, Unique, ...) a Taskok oldal tárgyalja, a borítékkal és a task-típus-elnevezéssel együtt, amik erre a kliensre épülnek.
Önálló példa
A queue-hoz semmi nem kell a kitből: egy kliens, egy task és egy minimális asynq szerver elég ahhoz, hogy végig lásd a munkát.
package main
import (
"context"
"log"
"github.com/gp-system/queue"
"github.com/hibiken/asynq"
)
func main() {
ctx := context.Background()
cfg := queue.Config{Addr: "localhost:6379"}
client := queue.MustNewClient(ctx, cfg)
defer client.Close()
task := asynq.NewTask("email:welcome", []byte(`{"user_id":"u_1"}`))
if _, err := client.Enqueue(ctx, task, queue.OnQueue("mail")); err != nil {
log.Fatal(err)
}
srv := asynq.NewServer(cfg.ValkeyConnOpt(), asynq.Config{Concurrency: 5})
mux := asynq.NewServeMux()
mux.HandleFunc("email:welcome", func(ctx context.Context, t *asynq.Task) error {
log.Printf("üdvözlő email küldése: %s", t.Payload())
return nil
})
if err := srv.Run(mux); err != nil {
log.Fatal(err)
}
}