Outbox
Importáld az outbox alcsomagot, hogy a dispatchered tranzakciós kézbesítési garanciát kapjon:
import "github.com/gp-system/events/outbox"
Az events/outbox a dispatch kézbesítési garanciája: az eventet nem közvetlenül Valkeybe adja, hanem egy Postgres-táblába írja a hívó tranzakciójában, és egy relay továbbítja a commitált sorokat a queue-ba. A generált projektben ez az alapértelmezett dispatcher: nem opcionális extra, hanem a kiindulópont.
A dual-write probléma
Vedd a shop PlaceOrder-ét: a rendelés insertje és az orderPlaced event együtt vagy sehogy kell megtörténjen. Két írás két rendszerbe (Postgres + Valkey) viszont nem tud atomikus lenni:
- Ha az eventet a commit előtt enqueue-lod, és a tranzakció rollbackel, a worker egy meg-nem-történt rendelésről küld visszaigazolót.
- Ha a commit után enqueue-lod, és a folyamat a kettő közt hal le (deploy, OOM, node-hiba), a rendelés megvan, de az event elveszett: a vevő soha nem kap e-mailt, és semmi nem jelzi, hogy baj van.
Önmagában egy commit-utáni hook nem zárja be ezt a rést: a rollback-esetet lekezelheti, de egy commit és a hook lefutása közti crash még mindig elveszítheti az eventet. A kitben ez a minta gyárilag jön: az outbox az eventet ugyanabba az adatbázisba, ugyanabban a tranzakcióban írja, mint az üzleti változást, egyetlen commit, ami vagy mindent perzisztál, vagy semmit.
Hogyan működik
A dispatch a tranzakcióban ír
A PlaceOrder a WithinTransaction-ön belül hívja a Dispatch-et; az outbox-dispatcher egy sort insertel az outbox_events táblába, ugyanazon a tranzakción át, mint a rendelés és a készlet írása.
Egyetlen commit
A tranzakció commitál: rendelés + készletcsökkentés + outbox-sor atomikusan. Rollbacknél az event-sor is eltűnik: fantomevent nincs.
A relay lefoglalja a commitált sorokat
A worker-ben futó relay OUTBOX_POLL_INTERVAL-onként pollozza a publikálatlan sorokat, FOR UPDATE SKIP LOCKED-dal, így N worker-replika soha nem verseng ugyanazokon a sorokon.
Enqueue determinisztikus TaskID-vel
Minden sort Valkeybe enqueue-l, asynq TaskID-nek az outbox event_id-t használva. Ha egy korábbi kör az enqueue után, de a megjelölés előtt halt le, az ismételt enqueue ErrDuplicate: a relay ezt sikernek veszi, és csak a megjelölést pótolja.
Megjelölés és takarítás
A beadott sorok published_at-et kapnak, és a foglaló tranzakció commitál. A publikált sorokat egy takarító ciklus OUTBOX_RETENTION után törli.
Innen az asynq veszi át: fan-out a listenerekre, retry, archiválás, lásd események. A lánc végig legalább-egyszeres, ezért a listener-idempotencia itt is előfeltétel.
Store és dispatcher
store := outbox.NewStore(pg.NewDB(pool)) // pg.DBTX-et fogad
dispatcher := outbox.NewDispatcher(store) // events.Dispatcher
A Store egyetlen metódus (Insert(ctx, env)), és a NewStore a kit pgx-executorára épül: mivel a pg.DB a ctx-ben hordozott tranzakcióhoz csatlakozik, a WithinTransaction-ön belüli insert automatikusan a hívó tranzakciójában fut: a service-nek semmit nem kell tudnia az outboxról. Tranzakción kívül hívva sima insert, az event így is kézbesítődik, csak az atomicitás-garancia nem értelmezett.
A bun-t használó projektek a events/outbox/bunx adaptert kötik be, ami a bun-tranzakcióhoz csatlakozik:
import outboxbunx "github.com/gp-system/events/outbox/bunx"
dispatcher := outbox.NewDispatcher(outboxbunx.NewStore(bunDB))
A generált main.go-kban ez a bekötés készen van (az első add event drótozza be); a shop HTTP-belépési pontjában így néz ki:
deps := shop.Dependencies{
DB: pg.NewDB(pool),
Transactor: pg.NewTransactor(pool),
Dispatcher: outbox.NewDispatcher(outbox.NewStore(pg.NewDB(pool))),
}
A relay
relay := outbox.MustNewRelay(ctx, cfg.Outbox, pool, cfg.Worker.Valkey)
err := worker.Run(ctx, cfg.Worker, register,
worker.WithOutboxRelay(relay), // a worker futtatja és állítja le
// ...
)
A NewRelay / MustNewRelay saját queue.Client-et nyit (PING-gel ellenőrizve), és a Run(ctx) a context lezárásáig polloz. Kézzel jellemzően egyiket sem hívod: a generált worker a worker.WithOutboxRelay-val adja át, és a worker életciklusa indítja/állítja le.
Amit a poll-ciklusról érdemes tudni:
- Tele batch → azonnali újra-poll. Ha egy kör
OUTBOX_BATCH_SIZE-nyi sort dolgozott fel, nem alszik, hanem azonnal újra polloz: a burst gyorsan leürül, aPOLL_INTERVALcsak az üresjárati késleltetés. - Hibánál exponenciális backoff. Egy hibázó kör (pl. Valkey-kiesés) duplázódó várakozással próbálkozik újra, 30 másodperces plafonig, majd az első sikeres kör visszaáll a normál ütemre.
- Részleges siker nem vész el. Ha a batch közepén hibázik az enqueue, a már beadott sorok megjelölődnek: a hiba nem játssza újra az egész batch-et.
Skálázhatóság: N replika, dupla-kézbesítés nélkül
A relay minden worker-replikában fut, koordináció nélkül. Két mechanizmus teszi ezt biztonságossá:
- a
FOR UPDATE SKIP LOCKEDmiatt két replika soha nem foglalja le ugyanazt a sort; - a determinisztikus
TaskID(=event_id) miatt a crash-ablakban kétszer beadott task a Valkeyben deduplikálódik: aOUTBOX_TASK_RETENTIONaz az időablak, ameddig a dedupe él.
Ezért skálázod a workert nyugodt szívvel vízszintesen, lásd worker → skálázás.
A tábla és a migráció
Az outbox_events goose-migrációként érkezik: a new project a migrations/20200101000100_outbox.sql-be írja, az add worker pedig ugyanezt a migrációt írja meg a worker-réteg bekötésének részeként; a séma forrása a csomagban az outbox.MigrationSQL. A lényege:
CREATE TABLE outbox_events (
id bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
event_id uuid NOT NULL UNIQUE,
event_name text NOT NULL,
payload jsonb NOT NULL,
metadata jsonb NOT NULL DEFAULT '{}'::jsonb,
created_at timestamptz NOT NULL DEFAULT now(),
published_at timestamptz
);
CREATE INDEX outbox_events_unpublished_idx ON outbox_events (id) WHERE published_at IS NULL;
A metadata a boríték trace-kontextusát hordozza, így a trace az outboxon át is folytonos; a parciális index a pollt tartja gyorsnak akkor is, ha a tábla történelmileg nagyra nő. Worker indítása előtt futtasd le: go run ./cmd/migrate up (lásd migrációk).
Konfiguráció
Az outbox.Config-ot a projekt configja OUTBOX_ prefixszel komponálja (a generált config.go-ban készen van):
| Változó | Default | Jelentés |
|---|---|---|
OUTBOX_POLL_INTERVAL | 1s | ennyit alszik a relay, ha nem tele batch-et kapott |
OUTBOX_BATCH_SIZE | 100 | egy poll legfeljebb ennyi sort foglal le és publikál |
OUTBOX_RETENTION | 168h | a publikált sorok ennyi ideig maradnak meg törlés előtt |
OUTBOX_CLEANUP_INTERVAL | 1h | ilyen gyakran fut a publikált sorok takarítása |
OUTBOX_TASK_RETENTION | 24h | a publikált taskok asynq-retenciója, egyben a TaskID-dedupe ablaka |
A kittel
A kiten kívül a migráció lefuttatása és a relay indítása a te felelősséged: az outbox.MigrationSQL-t azzal a migrációs eszközzel alkalmazod, amit a projekted amúgy is használ (a kit saját választása a goose, a cmd/migrate-en át), az outbox.Relay-t pedig magad indítod, a saját asynq szervered mellett. Egy generált projektben a new project/add worker megírja a migrációs fájlt, a worker chassis pedig végig a relay életciklusát birtokolja: neked csak az outbox.MustNewRelay-t kell hívnod, és átadnod a worker.WithOutboxRelay-nek.
Használt patternek
Ez az oldal a Design patternek katalógus Tranzakciós outbox bejegyzésének teljes kifejtése: a FOR UPDATE SKIP LOCKED batch-foglalás és a dual-write probléma itt, kódrészlettel és a shop PlaceOrder-példájával. Kanonikus külső leírás: microservices.io: Transactional outbox.