Wzorzec Transakcyjnego Skrzynki Wyjściowej w Go z PostgreSQL

Zapisz zdarzenie wraz z danymi. Nigdy ich nie rozdzielaj.

Page content

Dwa operacje zapisu, które powinny zostać wykonane wspólnie, w końcu niepowodzą się osobno.

Twój serwis zamówień zapisuje zamówienie w bazie danych, a następnie publikuje zdarzenie order.created w brokerze wiadomości.

Te dwie operacje wykonują się jedna po drugiej.

Pomiędzy nimi dochodzi do błędu: broker jest niedostępny, sieć się zawiesza, proces się restartuje lub kontener jest usuwany. Zapis do bazy danych powiódł się. Publikacja nie powiodła się. Serwis downstreamowy, który musi zostać poinformowany o nowym zamówieniu, nigdy się o tym nie dowiaduje. Nikt nie zauważył problemu, dopóki klient nie zadzwonił.

To jest problem podwójnego zapisu (dual-write) i jest jednym z najczęstszych źródeł cichej utraty danych w systemach rozproszonych. Standardowym rozwiązaniem jest wzorzec transakcyjnej skrzyni nadawczej (transactional outbox).

Transactional outbox pattern – event and data written together

Problem podwójnego zapisu (dual-write)

Tryb awarii jest prosty do zrozumienia, gdy się go zobaczy:

BEGIN;
  INSERT INTO orders ...   -- powiodło się
COMMIT;

PUBLISH order.created ...  -- niepowodzenie, awaria lub nigdy nie zostało wykonane

Baza danych i broker wiadomości nie dzielą wspólnego zakresu transakcji. Nie ma rollbacka, który obejmowałby obie operacje. Każda usługa wykonująca sekwencję save -> publish ma tę lukę. Wzorzec ten występuje w wielu postaciach:

  • db.Save(order) po którym następuje events.Publish(OrderCreated{...})
  • Handler HTTP, który zatwierdza transakcję, a następnie wywołuje zewnętrzny webhook
  • Worker, który przetwarza rekord z jednej kolejki i zapisuje wyniki do innej

Wynik we wszystkich przypadkach jest taki sam: jedna strona się powiodła, podczas gdy druga nie powiodła się, a system znajduje się w stanie niewidocznym dla monitoringu, ponieważ obie indywidualne operacje w pewnym momencie zwróciły sukces.

Pętla ponowionych prób (retry loop) nie rozwiązuje tego problemu. Ponawianie publikacji po zatwierdzeniu transakcji w bazie danych zadziała tylko wtedy, jeśli samo ponowienie będzie niezawodne – co wymaga gwarancji trwałości, której nie posiadasz.

Czym jest wzorzec transakcyjnej skrzyni nadawczej (transactional outbox)

Wzorzec outbox eliminuje lukę, całkowicie usuwając bezpośrednią publikację. Zamiast wywoływać brokera wewnątrz swojej logiki biznesowej, zapisujesz rekord zdarzenia do tabeli outbox w tej samej transakcji bazodanowej, co dane biznesowe. Oddzielny proces tła – przekaźnik (relay) – odczytuje z tabeli outbox i publikuje w brokerze.

BEGIN;
  INSERT INTO orders ...         -- dane biznesowe
  INSERT INTO outbox_events ...  -- rekord zdarzenia
COMMIT;

-- Proces przekaźnika (osobno):
SELECT ... FROM outbox_events FOR UPDATE SKIP LOCKED;
PUBLISH order.created ...
UPDATE outbox_events SET processed_at = NOW() WHERE id = $1;

Oba zapisy się powodzą lub obie nie powiodą się. Gwarancja transakcji, którą już posiadasz z PostgreSQL, obejmuje teraz również rekord zdarzenia. Przekaznik może ponawiać próbę publikowania tyle razy, ile to konieczne, ponieważ zdarzenie znajduje się w trwałej pamięci. Jeśli przekaźnik ulegnie awarii w trakcie działania, restartuje się i ponawia próbę. Najgorszym wynikiem jest opublikowanie zdarzenia więcej niż raz – co jest obsługiwane poprzez zapewnienie identyczności (idempotency) konsumentów (zobacz Idempotentność w systemach rozproszonych).

Schemat PostgreSQL dla tabeli outbox

Schemat jest celowo prosty:

CREATE TABLE outbox_events (
    id             UUID         PRIMARY KEY DEFAULT gen_random_uuid(),
    aggregate_type VARCHAR(100) NOT NULL,
    aggregate_id   VARCHAR(100) NOT NULL,
    event_type     VARCHAR(100) NOT NULL,
    payload        JSONB        NOT NULL,
    attempts       INT          NOT NULL DEFAULT 0,
    created_at     TIMESTAMPTZ  NOT NULL DEFAULT NOW(),
    processed_at   TIMESTAMPTZ
);

-- Indeks częściowy: indeksuje tylko nieprzetworzone wiersze, pozostaje mały po oznaczeniu wierszy jako przetworzonych
CREATE INDEX idx_outbox_unprocessed
    ON outbox_events (created_at)
    WHERE processed_at IS NULL;

Indeks częściowy na created_at WHERE processed_at IS NULL jest ważny. Bez niego indeks rośnie z każdym zapisanym zdarzeniem, a zapytanie wysyłające zapytania do przekaźnika (polling query) staje się wolniejsze z czasem. Dzięki niemu indeks obejmuje tylko oczekujące wiersze, które w stanie stałym są małym, ograniczonym zbiorem, niezależnie od liczby opublikowanych zdarzeń.

Kluczowe wybory pól:

  • aggregate_type i aggregate_id opisują, do której encji należy zdarzenie. Przydatne do gwarancji porządku i routingu.
  • event_type to nazwa zdarzenia, jakiej oczekują konsumenci.
  • payload JSONB przechowuje ciało zdarzenia. Użyj JSONB zamiast TEXT, aby móc go zapytać, jeśli zajdzie taka potrzeba.
  • attempts śledzi, ile razy przekaźnik próbował opublikować ten wiersz. Używane do limitów ponowionych prób i obsługi wiadomości trucicielskich (poison messages).
  • processed_at jest NULL dla oczekujących wierszy i ustawiane, gdy przekaźnik pomyślnie opublikuje zdarzenie.

Zapis danych biznesowych i zdarzenia outbox w jednej transakcji

Logika biznesowa zapisuje oba rekordy wewnątrz jednego wywołania BeginTx / Commit. Nie ma tu wywołania publikacji – tylko zapisy do bazy danych.

type OrderService struct {
    db *sql.DB
}

func (s *OrderService) CreateOrder(ctx context.Context, order Order) error {
    tx, err := s.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    if _, err := tx.ExecContext(ctx, `
        INSERT INTO orders (id, customer_id, total, created_at)
        VALUES ($1, $2, $3, NOW())
    `, order.ID, order.CustomerID, order.Total); err != nil {
        return fmt.Errorf("insert order: %w", err)
    }

    payload, err := json.Marshal(map[string]any{
        "order_id":    order.ID,
        "customer_id": order.CustomerID,
        "total":       order.Total,
    })
    if err != nil {
        return fmt.Errorf("marshal payload: %w", err)
    }

    if _, err := tx.ExecContext(ctx, `
        INSERT INTO outbox_events
            (aggregate_type, aggregate_id, event_type, payload)
        VALUES ($1, $2, $3, $4)
    `, "order", order.ID, "order.created", payload); err != nil {
        return fmt.Errorf("insert outbox event: %w", err)
    }

    return tx.Commit()
}

Jeśli tx.Commit() się nie powiedzie, ani wiersz zamówienia, ani wiersz outbox nie zostaną zapisane. Jeśli się powiedzie, oba są gwarantowanie obecne w bazie danych. Przekaznik może opublikować zdarzenie w dowolnym momencie po tym – natychmiast, po sekundzie lub po restarcie przekaźnika po awarii.

To jest jedyna zmiana kodu wymagana w warstwie biznesowej. Reszta wzorca znajduje się w przekaźniku.

Implementacja przekaźnika w Go

Przekaznik to worker tła, który regularnie odpytuje tabelę outbox. Pobiera partię nieprzetworzonych wierszy, publikuje każdy z nich i oznacza jako przetworzony. Możesz go utrzymywać w tym samym binarnym pliku co aplikację lub uruchomić jako osobny proces – oba podejścia działają, ale ten sam binarny plik jest prostszy w obsłudze.

type OutboxRelay struct {
    db          *sql.DB
    publisher   Publisher
    logger      *slog.Logger
    batchSize   int
    pollInterval time.Duration
    maxAttempts  int
}

func (r *OutboxRelay) Run(ctx context.Context) error {
    ticker := time.NewTicker(r.pollInterval)
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-ticker.C:
            if err := r.processBatch(ctx); err != nil {
                r.logger.Error("outbox relay batch failed", "err", err)
            }
        }
    }
}

Przekaznik szanuje anulowanie kontekstu, co ułatwia integrację z graceful shutdown. Szczegółowe omówienie żywotności kontekstu i wzorców anulowania znajdziesz w artykule Go context.Context Done Right.

FOR UPDATE SKIP LOCKED: wzorzec współbieżnego workera

Funkcja processBatch używa FOR UPDATE SKIP LOCKED, aby bezpiecznie obsługiwać współbieżnych workerów przekaźnika:

func (r *OutboxRelay) processBatch(ctx context.Context) error {
    tx, err := r.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    rows, err := tx.QueryContext(ctx, `
        SELECT id, aggregate_type, aggregate_id, event_type, payload
        FROM outbox_events
        WHERE processed_at IS NULL
          AND attempts < $1
        ORDER BY created_at
        LIMIT $2
        FOR UPDATE SKIP LOCKED
    `, r.maxAttempts, r.batchSize)
    if err != nil {
        return fmt.Errorf("query outbox: %w", err)
    }
    defer rows.Close()

    type row struct {
        id            string
        aggregateType string
        aggregateID   string
        eventType     string
        payload       json.RawMessage
    }

    var batch []row
    for rows.Next() {
        var e row
        if err := rows.Scan(
            &e.id, &e.aggregateType, &e.aggregateID, &e.eventType, &e.payload,
        ); err != nil {
            return fmt.Errorf("scan row: %w", err)
        }
        batch = append(batch, e)
    }
    if err := rows.Err(); err != nil {
        return err
    }

    for _, e := range batch {
        if err := r.publisher.Publish(ctx, e.eventType, e.aggregateID, e.payload); err != nil {
            r.logger.Error("publish failed", "event_id", e.id, "err", err)
            if _, err := tx.ExecContext(ctx,
                `UPDATE outbox_events SET attempts = attempts + 1 WHERE id = $1`, e.id,
            ); err != nil {
                r.logger.Error("increment attempts failed", "event_id", e.id, "err", err)
            }
            continue
        }

        if _, err := tx.ExecContext(ctx,
            `UPDATE outbox_events SET processed_at = NOW() WHERE id = $1`, e.id,
        ); err != nil {
            return fmt.Errorf("mark processed: %w", err)
        }
    }

    return tx.Commit()
}

FOR UPDATE SKIP LOCKED wykonuje dwie rzeczy. Po pierwsze, FOR UPDATE blokuje wybrane wiersze przez czas trwania transakcji, zapobiegając innym transakcjom na ich selekcję. Po drugie, SKIP LOCKED oznacza, że jeśli wiersz jest już zablokowany przez inną transakcję, zapytanie pomija go zamiast czekać. Efektem jest to, że wiele workerów przekaźnika może działań równolegle, a każdy wybierze nieprzecinający się podzbiór wierszy.

Bez SKIP LOCKED drugi worker zablokowałby się do zatwierdzenia transakcji przez pierwszy worker przed zobaczeniem tych samych wierszy – w tym momencie byłyby one już oznaczone jako przetworzone. Dzięki SKIP LOCKED drugi worker natychmiast wybiera inne wiersze zamiast czekać, co daje bezpieczne skalowanie poziome.

Zwróć uwagę na separację skanowania i publikacji w kodzie powyżej: wszystkie wiersze są skanowane do tablicy przed rozpoczęciem pętli publikacji. Unikamy w ten sposób trzymania otwartego kursora *sql.Rows podczas wywołań sieciowych do brokera, co wydłużyłoby czas trzymania transakcji.

Idempotentność i deduplikacja

Przekaznik publikuje przynajmniej raz (at-least-once). Jeśli opublikuje zdarzenie, a następnie ulegnie awarii przed zatwierdzeniem aktualizacji processed_at, opublikuje to samo zdarzenie ponownie po restarcie. Jest to nieuniknione – dostawa dokładnie raz (exactly-once) między bazą danych a brokerem bez koordynatora transakcji rozproszonych wymaga tego kompromisu.

Konsumenci muszą być idempotentni. Najprostszym podejściem jest śledzenie przetworzonych identyfikatorów zdarzeń w tabeli processed_events:

CREATE TABLE processed_events (
    event_id   UUID PRIMARY KEY,
    processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
func (h *OrderHandler) HandleOrderCreated(ctx context.Context, eventID string, payload []byte) error {
    // Deduplikacja przy użyciu identyfikatora zdarzenia jako klucza naturalnego
    _, err := h.db.ExecContext(ctx, `
        INSERT INTO processed_events (event_id) VALUES ($1)
        ON CONFLICT (event_id) DO NOTHING
    `, eventID)
    if err != nil {
        return fmt.Errorf("dedup check: %w", err)
    }

    // Sprawdź, czy insert faktycznie się wydarzył (1 wiersz) czy był no-op (0 wierszy)
    // Prostszе podejście: użyj RETURNING lub sprawdź liczbę affected rows
    // Jeśli 0 wierszy affected, to jest duplikat -- pomiń go
    ...
}

W praktyce wiele zespołów polega na własnej deduplikacji brokera (takiej jak pole key w Kafka dla tematów z kompresją logów lub nagłówek message-id w RabbitMQ) i traktuje deduplikację na poziomie bazy danych jako zapasowe rozwiązanie. Oba podejścia są poprawnymi warstwami do zastosowania.

Dołącz identyfikator zdarzenia outbox id (UUID) do opublikowanej wiadomości jako klucz deduplikacji. Konsumenci mogą go wtedy użyć, niezależnie od preferowanego mechanizmu deduplikacji.

Polityka ponowionych prób i wiadomości trucicielskie

Kolumna attempts steruje polityką ponowionych prób. Przekaznik pomija wiersze, gdzie attempts >= maxAttempts i traktuje je jako wiadomości trucicielskie (dead letters). Oddzielny proces lub alert operatora zajmuje się nimi.

Prosty widok wiadomości trucicielskich:

CREATE VIEW outbox_dead_letters AS
SELECT *
FROM outbox_events
WHERE attempts >= 5
  AND processed_at IS NULL
ORDER BY created_at;

Dobra polityka ponowionych prób w produkcji:

  • Ustaw maxAttempts na 5-10 w zależności od kosztu ponowionych prób.
  • Rozważ wykładniczy backoff: dodaj kolumnę retry_after i pomijaj wiersze, gdzie retry_after > NOW().
  • Generuj alerty, gdy COUNT(*) FROM outbox_dead_letters przekroczy próg.
  • Zapewnij ręczną ścieżkę ponowienia: endpoint administratora lub skrypt, który resetuje attempts = 0 i retry_after = NULL dla konkretnych wierszy.

Wiadomości trucicielskie – wiersze, które konsekwentnie się nie powiodą z powodu błędu w konsumentcie lub niezgodności schematu – nie powinny blokować zdrowych wiadomości. Ponieważ przekaźnik przetwarza partię na tick i oznacza niepowodzenia inkrementacją prób zamiast usuwając je z kolejki, zdrowe wiersze postępują normalnie, podczas gdy trucicielskie gromadzą próby, aż osiągną próg wiadomości trucicielskich. Ten widok outbox_dead_letters jest wersją na poziomie bazy danych tego samego wzorca, który natywnie implementują kolejki wiadomości trucicielskich). – kwarantanna po przekroczeniu progu, alert przy dużej objętości i wymóg świadomej decyzji przed odtworzeniem.

Porządek zdarzeń i partycjonowanie

Zapytanie wysyłające zapytania do przekaźnika (polling query) sortuje po created_at, co daje porządek FIFO wewnątrz partii. Dla większości przypadków użycia to wystarczy. Gdy zależy nam na ścisłym porządku dla każdej encji – na przykład zapewnienie, że order.updated nigdy nie jest publikowane przed order.created dla tego samego zamówienia – potrzebujemy porządku per-agregat.

Dodaj aggregate_id do klauzuli ORDER BY i użyj go jako klucza wiadomości podczas publikowania do tematu partycjonowanego, takiego jak Apache Kafka. Kafka kieruje wszystkie wiadomości z tym samym kluczem do tej samej partycji, a partycje są konsumowane w porządku. Daje to gwarancje porządku per-agregat bez globalnego porządku, który wymagałby pojedynczej instancji przekaźnika.

ORDER BY aggregate_id, created_at

Dla brokerów, które nie obsługują porządku partycjonowanego (takich jak podstawowe kolejki AMQP), pojedyncza instancja przekaźnika lub sprawdzanie porządku na poziomie aplikacji w konsumentcie są praktycznymi alternatywami.

Zmniejszenie opóźnień pollingiem za pomocą LISTEN/NOTIFY

Interwał pollingowy jednej sekundy oznacza średnie opóźnienie zdarzeń rzędu 500 milisekund. Dla większości obciążeń to jest w porządku. Dla przypadków, gdzie potrzebujesz bliskiej zeru opóźnienia, mechanizm LISTEN/NOTIFY PostgreSQL pozwala przekaźnikowi obudzić się natychmiast po wstawieniu nowego wiersza outbox.

Dodaj trigger do tabeli outbox:

CREATE OR REPLACE FUNCTION notify_outbox_insert() RETURNS trigger AS $$
BEGIN
    PERFORM pg_notify('outbox_event', NEW.id::text);
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER outbox_insert_notify
AFTER INSERT ON outbox_events
FOR EACH ROW EXECUTE FUNCTION notify_outbox_insert();

W przekaźniku nasłuchuj na kanale i budź się na notyfikacje, nadal pozostawiając okresowy polling jako fallback:

func (r *OutboxRelay) Run(ctx context.Context) error {
    listener := pq.NewListener(r.dsn, 10*time.Second, time.Minute, nil)
    defer listener.Close()

    if err := listener.Listen("outbox_event"); err != nil {
        return fmt.Errorf("listen: %w", err)
    }

    ticker := time.NewTicker(5 * time.Second) // fallback poll
    defer ticker.Stop()

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        case <-listener.Notify:
            if err := r.processBatch(ctx); err != nil {
                r.logger.Error("outbox batch failed (notify)", "err", err)
            }
        case <-ticker.C:
            if err := r.processBatch(ctx); err != nil {
                r.logger.Error("outbox batch failed (poll)", "err", err)
            }
        }
    }
}

Fallback ticker obsługuje wszystkie pominięte notyfikacje podczas restartu przekaźnika lub migawki sieciowej. Utrzymuj interwał fallbackowy na kilku sekundach, a nie milisekundach – jego zadaniem jest odzyskiwanie, a nie niskie opóźnienia.

Obserwowalność: metryki, logi i alerty

Outbox to infrastruktura. Traktuj go jak infrastrukturę i odpowiednio go instrumentuj.

Kluczowe metryki:

var (
    outboxPublished = prometheus.NewCounter(prometheus.CounterOpts{
        Name: "outbox_events_published_total",
        Help: "Total outbox events successfully published.",
    })
    outboxFailed = prometheus.NewCounterVec(prometheus.CounterOpts{
        Name: "outbox_events_failed_total",
        Help: "Total outbox publish failures by event type.",
    }, []string{"event_type"})
    outboxPending = prometheus.NewGauge(prometheus.GaugeOpts{
        Name: "outbox_events_pending",
        Help: "Current number of unprocessed outbox events.",
    })
    outboxBatchDuration = prometheus.NewHistogram(prometheus.HistogramOpts{
        Name:    "outbox_batch_duration_seconds",
        Help:    "Duration of each outbox processing batch.",
        Buckets: prometheus.DefBuckets,
    })
)

Odświeżanie Gauge: uruchom okresowe zapytanie, aby utrzymać outbox_events_pending dokładnym:

SELECT COUNT(*) FROM outbox_events WHERE processed_at IS NULL;

Progi alertów do rozważenia:

  • outbox_events_pending > 1000 przez więcej niż dwie minuty: przekaźnik się zapycha lub utknął.
  • outbox_events_pending rośnie monotonicznie: broker jest niedostępny lub przekaźnik uległ awarii.
  • Liczba wiadomości trucicielskich niezerowa: błąd schematu lub konsumenta wymaga zbadania.
  • outbox_batch_duration_seconds p95 > 5s: baza danych jest wolna lub rozmiar partii jest zbyt duży.

Pola strukturalnych logów: dołącz event_id, event_type, aggregate_id i attempt do każdej linii logu z przekaźnika. Pola te pozwalają skorelować nieudaną publikację z konkretnym wierszem outbox i śledzeniem konsumenta downstream.

Outbox vs. bezpośrednia kolejka vs. Saga

Wzorzec outbox nie jest odpowiednim narzędziem do każdego problemu koordynacji. Oto porównanie:

Podejście Atomowość Złożoność Kiedy stosować
Bezpośrednia publikacja Brak Niska Akceptowalne przy okresowej utracie zdarzeń
Transakcyjny outbox Silna Średnia Niezawodna dostawa zdarzeń z pojedynczej usługi
Wzorzec Saga Ostateczna Wysoka Transakcje międzyusługowe obejmujące wiele baz danych
Dwa fazy komitów (Two-phase commit) Silna Bardzo wysoka Rzadko praktyczne; unikane w większości systemów rozproszonych

Wzorzec outbox gwarantuje, że pojedyncza usługa niezawodnie emituje zdarzenia, które odzwierciedlają jej własne zmiany stanu. Nie koordynuje zmian stanu między wieloma usługami – do tego służy Wzorzec Saga. Wybór brokera – czy to RabbitMQ, SQS, czy Kafka – jest niezależny od samego wzorca outbox; przekaźnik publikuje do dowolnego brokera, którego używa Twój system.

Jeśli budujesz Sagę, wzorzec outbox nadal jest użyteczny: każdy uczestnik Sagii zapisuje swoją lokalną zmianę stanu i swoje zdarzenie Sagii w jednej transakcji za pomocą outbox, a następnie orchestrator lub choreography Sagii odczytuje te zdarzenia niezawodnie.

CDC na bazie WAL jako alternatywny przekaźnik

Zamiast pollingowania, możesz śledzić Write-Ahead Log (WAL) PostgreSQL i czytać wstawienia outbox bezpośrednio ze strumienia replikacji. Narzędzia takie jak Debezium to robią. Zalety to niższe opóźnienia i brak presji blokad na tabeli outbox. Wady to złożoność operacyjna, dedykowany slot replikacji PostgreSQL i zewnętrzny serwis do uruchamiania i monitorowania.

Dla większości zespołów przekaźnik opisany powyżej jest dobrym punktem startowym. Śledzenie WAL ma sens, gdy masz wysokie wsady wstawień do outbox (dziesiątki tysięcy na sekundę), potrzebujesz opóźnień zdarzeń poniżej 100ms lub już uruchamiasz Debezium dla innych potrzeb przechwytywania zmian.

Integracja z sqlc

Jeśli używasz sqlc do typowo bezpiecznego kodu Go bazy danych, zapytania outbox wpisują się naturalnie:

-- name: InsertOutboxEvent :exec
INSERT INTO outbox_events (aggregate_type, aggregate_id, event_type, payload)
VALUES (@aggregate_type, @aggregate_id, @event_type, @payload);

-- name: FetchOutboxBatch :many
SELECT id, aggregate_type, aggregate_id, event_type, payload
FROM outbox_events
WHERE processed_at IS NULL
  AND attempts < @max_attempts
ORDER BY created_at
LIMIT @batch_size
FOR UPDATE SKIP LOCKED;

-- name: MarkOutboxProcessed :exec
UPDATE outbox_events SET processed_at = NOW() WHERE id = @id;

-- name: IncrementOutboxAttempts :exec
UPDATE outbox_events SET attempts = attempts + 1 WHERE id = @id;

-- name: OutboxPendingCount :one
SELECT COUNT(*) FROM outbox_events WHERE processed_at IS NULL;

sqlc generuje funkcje typu bezpiecznego dla każdego zapytania, co unika błędów interpolacji ciągów i utrzymuje logikę zapytań outbox współlokalizowaną z resztą warstwy dostępu do bazy danych.

Lista kontrolna produkcji

Użyj tego przed wypuszczeniem implementacji outbox:

Baza danych

  • Tabela outbox ma indeks częściowy na created_at WHERE processed_at IS NULL
  • Kolumna attempts obecna z domyślną wartością 0
  • Zdefiniowany widok wiadomości trucicielskich lub zapytanie
  • Stare przetworzone wiersze są okresowo archiwizowane lub usuwane (wystarczy zadanie czyszczenia w nocy)

Przekaznik

  • FOR UPDATE SKIP LOCKED użyte w zapytaniu pollingowym
  • Przekaznik działa wewnątrz transakcji (begin przed zapytaniem, commit po wszystkich aktualizacjach)
  • Rozmiar partii jest ograniczony (50-200 wierszy jest typowe)
  • Przekaznik szanuje anulowanie kontekstu dla graceful shutdown
  • Nieudane publikacje inkrementują attempts zamiast powodować zatrzymanie partii

Idempotentność

Obserwowalność

  • Gauge outbox_events_pending jest monitorowany i alertowany
  • Liczba wiadomości trucicielskich jest alertowana
  • Czas trwania partii przekaźnika jest śledzony
  • Strukturalne logi zawierają event_id, event_type i aggregate_id

Operacje

  • Istnieje ścieżka ręcznego ponowienia dla wierszy wiadomości trucicielskich
  • Zachowanie restartu przekaźnika jest testowane (czy publikuje poprawnie ponownie?)
  • Zachowanie awarii brokera jest testowane (czy outbox rośnie i opróżnia się poprawnie?)

Zakończenie

Problem podwójnego zapisu jest łatwy do odrzucenia jako przypadek brzegowy, dopóki nie spowoduje incydentu. Wzorzec transakcyjnej skrzyni nadawczej rozwiązuje go za pomocą narzędzi, które już posiadasz: transakcji PostgreSQL, goroutiny tła i jednej dodatkowej tabeli. Przekaznik jest prosty w budowie, prosty w obsłudze i prosty do zrozumienia.

Koszt polega na tym, że konsumenci muszą być zaprojektowani dla dostawy przynajmniej raz (at-least-once). To jest rozsądny kompromis. Dostawa dokładnie raz (exactly-once) między bazą danych a brokerem bez transakcji rozproszonych nie jest praktycznie możliwa – a udawanie, że tak jest, prowadzi do systemów, które cicho porzucają lub podwójnie przetwarzają zdarzenia w warunkach awarii.

Zapisz zdarzenie wraz z danymi. Przekaż je niezawodnie. Zrób konsumentów idempotentnymi. To jest cały wzorzec.

Ten artykuł jest częścią klastra Architektura aplikacji w produkcji.

Źródła

Subskrybuj

Otrzymuj nowe wpisy o systemach, infrastrukturze i inżynierii AI.