events

Outbox

A tranzakciós outbox minta, az üzleti írás és az event atomikusan, kézbesítési garanciával.

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:

cmd/shop/main.go (részlet)
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, a POLL_INTERVAL csak 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á:

  1. a FOR UPDATE SKIP LOCKED miatt két replika soha nem foglalja le ugyanazt a sort;
  2. a determinisztikus TaskID (= event_id) miatt a crash-ablakban kétszer beadott task a Valkeyben deduplikálódik: a OUTBOX_TASK_RETENTION az 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óDefaultJelentés
OUTBOX_POLL_INTERVAL1sennyit alszik a relay, ha nem tele batch-et kapott
OUTBOX_BATCH_SIZE100egy poll legfeljebb ennyi sort foglal le és publikál
OUTBOX_RETENTION168ha publikált sorok ennyi ideig maradnak meg törlés előtt
OUTBOX_CLEANUP_INTERVAL1hilyen gyakran fut a publikált sorok takarítása
OUTBOX_TASK_RETENTION24ha 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.

Copyright © 2026