Ütemezés
Importáld a scheduler alcsomagot, hogy cron-ütemezéseket és jobokat kódban definiálj:
import "github.com/gp-system/events/scheduler"
A scheduler lehetővé teszi, hogy az ütemezésed kódban definiáld, ne crontabban: egy Schedule-t építesz fel fluensen, és a worker futtatja. Az, hogy a replikák közül csak egyen fut, itt alapértelmezett viselkedés, nem opt-in: N worker-replika közül Valkey-lease leader-election dönti el, melyik tüzel, külön flag, külön cache-lock nélkül.
Két dolgot ütemezhetsz:
- eventet: a tüzelés egy sima dispatch, az event szétterül az összes listenerére;
- jobot: egyetlen nevesített handler, fan-out nélkül.
A schedule felépítése
Minden modul generált jobs/register.go-ja ad egy Register(sched, deps)-t; a worker ezeket egyetlen schedule-be komponálja. A shop éjszakai riportját az add job shop nightlySalesReport --cron "0 3 * * *" generálta:
package jobs
import (
"github.com/acme/shop/internal/modules/shop"
"github.com/gp-system/gpsystem/scheduler"
// gpsystem:job-imports
)
func Register(sched *scheduler.Schedule, deps shop.Dependencies) {
sched.Job("shop.nightlySalesReport", "0 3 * * *", nightlySalesReport(deps))
// gpsystem:jobs
}
A job handler-stubja külön fájlban, a törzs a tiéd:
package jobs
import (
"context"
"github.com/acme/shop/internal/modules/shop"
)
func nightlySalesReport(deps shop.Dependencies) func(context.Context) error {
return func(ctx context.Context) error {
report, err := buildSalesReport(ctx, deps.DB)
if err != nil {
return err
}
return deps.Mailer.Send(ctx, salesReportMessage(report))
}
}
A visszaadott hiba retry-t vált ki (az alapértelmezett 25-ös keretig, vagy a bejegyzés saját queue.MaxRetry-áig), egy átmeneti DB-hiba miatt nem marad el a riport.
Event-schedule vagy job-schedule
A Schedule fluent API-ja: minden metódus *Schedule-t ad vissza, láncolható.
| Metódus | Tüzel |
|---|---|
Cron(spec, ev, opts...) | cron-kifejezés szerint |
Every(d, ev, opts...) | d időközönként (@every descriptorra fordul) |
EveryMinute(ev, opts...) | minden perc elején |
Hourly(ev, opts...) | minden óra elején |
Daily(ev, opts...) | éjfélkor, a SCHEDULER_TIMEZONE szerint |
DailyAt("HH:MM", ev, opts...) | naponta a megadott időpontban (24 órás; hibás formátumra setupkor panikol) |
Job(name, spec, fn, opts...) | cron szerint futtatja az fn func(context.Context) error-t |
Mikor melyik? Event-schedule, ha a tüzelésre több, egymástól függetlenül retry-zott reakció kell, vagy ha ugyanarra az eseményre futásidejű dispatch is létezik: a listenerek nem tudják és nem is kell tudniuk, hogy cron vagy user-akció dobta el az eventet. Job-schedule, ha egyetlen, nevesített teendőről van szó, mint a shop riportja.
// Event-schedule: minden tüzelés friss eventet dispatchel, ami szétterül
// az összes listenerére, pont mint egy futásidejű Dispatch.
sched.Every(15*time.Minute, shopevents.SyncInventory{})
sched.DailyAt("06:30", shopevents.DigestDue{})
// Job-schedule: egy handler, fan-out nélkül; a task-típus "job:<név>".
sched.Job("shop.nightlySalesReport", "0 3 * * *", nightlySalesReport(deps))
Az opts... a szokásos enqueue-opciók (pl. queue.OnQueue("reports"), queue.Timeout(10*time.Minute)), és a tüzeléskor beadott taskra vonatkoznak.
Egy technikai részlet, ami a garanciákat adja: az ütemezett event tüzelésekor a worker a boríték id-jának a trigger-task id-ját használja. Így minden tüzelés új dispatch (minden listener újra megkapja), de ugyanannak a tüzelésnek a retry-ja idempotens: nem duplázza a fan-outot.
Cron-szintaxis
A spec a standard 5 mezős cron (perc óra hónap-napja hónap hét-napja), plusz az asynq descriptorai:
"0 3 * * *" minden nap 03:00-kor
"*/10 * * * *" 10 percenként
"0 8 * * 1" hétfőnként 08:00-kor
"@every 30m" 30 percenként (nem az óra rácsán, az indulástól számít)
"@daily" éjfélkor
A kifejezések a SCHEDULER_TIMEZONE-ban értékelődnek ki (IANA név, pl. Europe/Budapest); a default UTC, tehát a shop "0 3 * * *"-a UTC 03:00.
Replika-biztonság: a leader-election
A schedule-t a Runner hajtja végre, amit sosem konstruálsz kézzel: a worker.Run automatikusan elindítja, ha a schedule nem üres. Belül az asynq Scheduler-e dolgozik, egy Valkey-lease védelmében:
- Minden worker-replika megpróbálja megszerezni a lease-kulcsot (
SETNX,SCHEDULER_LEASE_TTLlejárattal). - Csak a lease birtokosa regisztrálja és futtatja a cron-bejegyzéseket;
LEASE_TTL/3-onként megújítja a lease-t (Lua compare-and-act, hogy egy lemaradt régi leader ne írhassa felül az újat). - Ha a leader lehal vagy elveszti a lease-t, leáll a tüzelése, és legfeljebb
LEASE_TTLmúlva egy másik replika veszi át.
Átvételkor rövid átfedés elképzelhető (a régi leader még kiadott egy ticket, az új is), de ezt a rendszer többi része nyeli el: a kézbesítés amúgy is legalább-egyszeres, a listenerek idempotensek. Ha több alkalmazás (vagy környezet) osztozik egy Valkeyen, a SCHEDULER_NAMESPACE-szel válaszd szét a lease-kulcsaikat, különben egymás leader-választásán osztoznak.
// az "egy szerveren fusson" itt nem opció, hanem az egyetlen működési mód
sched.Job("shop.nightlySalesReport", "0 3 * * *", nightlySalesReport(deps))
Mi fut hol
Fontos különbség a crontabhoz képest: a tüzelés és a végrehajtás szétválik. A leader csak enqueue-l: a schedule:<event> vagy job:<név> task a queue-ba kerül, és bármelyik worker-replika feldolgozhatja, a szokásos konkurencia- és queue-súlyok szerint. Egy nehéz riport tehát nem a "scheduler gépet" terheli, hanem a poolt.
Konfiguráció
A scheduler.Config a worker.Config-ba ágyazva érkezik, SCHEDULER_ prefixszel:
| Változó | Default | Jelentés |
|---|---|---|
SCHEDULER_TIMEZONE | UTC | IANA időzóna, amelyben a cron-kifejezések kiértékelődnek |
SCHEDULER_LEASE_TTL | 15s | a leader-lease TTL-je; a leader TTL/3-onként újítja meg |
SCHEDULER_NAMESPACE | gpsystem | a lease-kulcs névtere; állítsd az app nevére, ha többen osztoztok egy Valkeyen |
SCHEDULER_NAMESPACE alapértéke szó szerint a gpsystem string, nem a projekted neve. Ez egy generált kit-projekten belül ártalmatlan (minden projektnek saját Valkey-instance-a van), de ha az events/scheduler-t önállóan használod, és több, egymástól független service-t ugyanarra a Valkeyre mutatsz, mindegyik ugyanazt a lease-kulcsot próbálja megszerezni, hacsak explicit be nem állítod. Állítsd az app nevére.A kittel
A worker chassis az, ami ténylegesen felépíti és futtatja a Runner-t: a worker.Run automatikusan elindítja, amint a RegisterFunc-od által felépített Schedule nem üres, és ugyanabban a leállási sorrendben állítja le, mint az asynq szervert és az outbox relay-t. Ebben semmi kit-specifikus nincs: egy önálló fogyasztó felépít egy Schedule-t, megnyit egy queue.Config-ot, és közvetlenül hívja a scheduler.NewRunner(ctx, cfg, queueCfg, sched).Run(ctx)-t.
Használt patternek
Valkey-lease leader election (SETNX szerzéshez, Lua compare-and-renew a megújításhoz): lásd a Design patterneket.