events

Ütemezés

Task scheduling kódban, cron-eventek és jobok, replika-biztosan, Valkey-lease leader-electionnel.

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:

internal/modules/shop/jobs/register.go
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:

internal/modules/shop/jobs/nightly_sales_report.go
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ódusTü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:

  1. Minden worker-replika megpróbálja megszerezni a lease-kulcsot (SETNX, SCHEDULER_LEASE_TTL lejárattal).
  2. 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).
  3. Ha a leader lehal vagy elveszti a lease-t, leáll a tüzelése, és legfeljebb LEASE_TTL mú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.

gpsystem: internal/modules/shop/jobs/register.go
// 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óDefaultJelentés
SCHEDULER_TIMEZONEUTCIANA időzóna, amelyben a cron-kifejezések kiértékelődnek
SCHEDULER_LEASE_TTL15sa leader-lease TTL-je; a leader TTL/3-onként újítja meg
SCHEDULER_NAMESPACEgpsystema lease-kulcs névtere; állítsd az app nevére, ha többen osztoztok egy Valkeyen
A 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.

Copyright © 2026