Le modèle Outbox transactionnel en Go avec PostgreSQL

Écrivez l'événement avec les données. Ne les séparez jamais.

Sommaire

Deux écritures qui devraient réussir ensemble finiront par échouer séparément. Votre service de commande enregistre la commande dans la base de données, puis publie un événement order.created vers un courtier de messages (message broker).

Ces deux opérations s’exécutent l’une après l’autre.

Entre les deux, des choses peuvent mal tourner : le courtier est indisponible, le réseau expire, le processus redémarre ou le conteneur est évacué. L’écriture dans la base de données a réussi. La publication n’a pas abouti. Le service en aval qui doit être informé de la nouvelle commande ne l’apprend jamais. Personne ne s’en rend compte jusqu’à ce qu’un client appelle.

C’est le problème du double écriture (dual-write), et il est l’une des sources les plus courantes de perte silencieuse de données dans les systèmes distribués. Le modèle de la boîte aux lettres transactionnelle (transactional outbox) est la correction standard.

Schéma du modèle de la boîte aux lettres transactionnelle – événement et données écrits ensemble

Le problème du double écriture

Le mode de défaillance est facile à comprendre une fois que vous l’avez identifié :

BEGIN;
  INSERT INTO orders ...   -- réussit
COMMIT;

PUBLISH order.created ...  -- échoue, plante ou n'est jamais atteint

La base de données et le courtier de messages ne partagent pas de limite de transaction. Il n’y a pas de retour en arrière (rollback) qui couvre les deux. Chaque service qui effectue save -> publish séquentiellement présente cette lacune. Le modèle se présente sous plusieurs formes :

  • db.Save(order) suivi de events.Publish(OrderCreated{...})
  • Gestionnaire HTTP qui valide une transaction puis appelle un webhook externe
  • Worker qui traite un enregistrement d’une file d’attente et écrit les résultats dans une autre

L’issue est la même dans tous les cas : un côté réussit pendant que l’autre échoue, et le système se retrouve dans un état invisible pour la supervision car les deux opérations individuelles ont renvoyé un succès à un moment donné.

Une boucle de nouvelle tentative ne résout pas ce problème. Relancer la publication après la validation de la base de données ne fonctionne que si la nouvelle tentative elle-même est fiable – ce qui nécessite la garantie de durabilité exacte que vous n’avez pas.

Ce que fait le modèle de la boîte aux lettres transactionnelle

Le modèle outbox élimine l’écart en supprimant complètement la publication directe. Au lieu d’appeler le courtier depuis votre logique métier, vous écrivez un enregistrement d’événement dans une table outbox lors de la même transaction de base de données que les données métier. Un processus distinct – le relais (relay) – lit depuis la table outbox et publie vers le courtier.

BEGIN;
  INSERT INTO orders ...         -- données métier
  INSERT INTO outbox_events ...  -- enregistrement d'événement
COMMIT;

-- Processus Relay (séparément) :
SELECT ... FROM outbox_events FOR UPDATE SKIP LOCKED;
PUBLISH order.created ...
UPDATE outbox_events SET processed_at = NOW() WHERE id = $1;

Soit les deux écritures réussissent, soit elles échouent. La garantie de transaction que vous avez déjà de PostgreSQL couvre désormais également l’enregistrement d’événement. Le relais peut relancer la publication autant de fois que nécessaire car l’événement repose dans un stockage durable. Si le relais plante en cours de route, il redémarre et réessaie. Le pire scénario est que l’événement soit publié plus d’une fois – ce qui est géré en rendant les consommateurs idempotents (voir Idempotence dans les Systèmes Distribués).

Schéma PostgreSQL pour la table outbox

Le schéma est intentionnellement simple :

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
);

-- Index partiel : indexe uniquement les lignes non traitées, reste petit au fur et à mesure que les lignes sont marquées comme terminées
CREATE INDEX idx_outbox_unprocessed
    ON outbox_events (created_at)
    WHERE processed_at IS NULL;

L’index partiel sur created_at WHERE processed_at IS NULL est important. Sans lui, l’index grandit avec chaque événement jamais écrit et la requête de sondage du relais devient plus lente avec le temps. Avec lui, l’index ne couvre que les lignes en attente, qui dans un état stable constituent un petit ensemble borné, indépendamment du nombre d’événements publiés.

Choix clés des champs :

  • aggregate_type et aggregate_id décrivent à quelle entité appartient l’événement. Utile pour les garanties de tri et le routage.
  • event_type est le nom de l’événement que vos consommateurs attendent.
  • payload JSONB stocke le corps de l’événement. Utilisez JSONB plutôt que TEXT afin de pouvoir l’interroger si nécessaire.
  • attempts suit le nombre de fois où le relais a essayé de publier cette ligne. Utilisé pour les limites de nouvelle tentative et la gestion des messages empoisonnés.
  • processed_at est NULL pour les lignes en attente et défini lorsque le relais publie avec succès.

Écriture des données métier et de l’événement outbox dans une seule transaction

La logique métier écrit les deux enregistrements à l’intérieur d’un seul appel BeginTx / Commit. Il n’y a aucun appel de publication ici – uniquement des écritures de base de données.

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()
}

Si tx.Commit() échoue, ni la ligne de commande ni la ligne outbox ne sont persistées. Si elle réussit, toutes deux sont garanties d’être dans la base de données. Le relais peut publier l’événement à tout moment après cela – immédiatement, dans une seconde, ou après le redémarrage du relais suite à un plantage.

C’est le seul changement de code requis dans votre couche métier. Le reste du modèle réside dans le relais.

Implémentation du relais Go

Le relais est un travailleur en arrière-plan qui effectue un sondage (polling) de la table outbox selon un temporisateur. Il récupère un lot de lignes non traitées, publie chacune d’elles, puis les marque comme terminées. Gardez-le dans le même binaire que votre application ou exécutez-le en tant que processus distinct – les deux fonctionnent, mais le même binaire est plus simple à gérer.

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)
            }
        }
    }
}

Le relais respecte l’annulation du contexte, ce qui facilite son intégration avec l’arrêt progressif (graceful shutdown). Pour un traitement détaillé de la durée de vie du contexte et des modèles d’annulation, consultez Go context.Context Done Right.

FOR UPDATE SKIP LOCKED : le modèle de travailleur concurrent

La fonction processBatch utilise FOR UPDATE SKIP LOCKED pour gérer en toute sécurité les travailleurs relais concurrents :

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 fait deux choses. Premièrement, FOR UPDATE verrouille les lignes sélectionnées pendant toute la durée de la transaction, empêchant toute autre transaction de les sélectionner. Deuxièmement, SKIP LOCKED signifie que si une ligne est déjà verrouillée par une autre transaction, la requête la saute au lieu d’attendre. Le résultat est que plusieurs travailleurs relais peuvent s’exécuter en parallèle et chacun récupérera un sous-ensemble non chevauchant de lignes.

Sans SKIP LOCKED, un deuxième travailleur bloquerait jusqu’à ce que la première transaction soit validée avant de voir les mêmes lignes – auquel point elles seraient déjà marquées comme terminées. Avec SKIP LOCKED, le deuxième travailleur récupère immédiatement des lignes différentes au lieu d’attendre, ce qui vous offre une mise à l’échelle horizontale sûre.

Notez la séparation entre lecture et publication dans le code ci-dessus : toutes les lignes sont lues dans un tableau avant le début de la boucle de publication. Cela évite de maintenir un curseur *sql.Rows ouvert pendant les appels réseau vers le courtier, ce qui maintiendrait la transaction ouverte plus longtemps que nécessaire.

Idempotence et déduplication

Le relais publie au moins une fois. S’il publie un événement puis plante avant de valider la mise à jour de processed_at, il publiera le même événement encore au redémarrage. Cela est inévitable – une livraison exactement une fois (exactly-once) entre une base de données et un courtier de messages sans coordinateur de transactions distribuées exige ce compromis.

Les consommateurs doivent être idempotents. L’approche la plus simple consiste à suivre les ID d’événements traités dans une table 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 {
    // Déduplication en utilisant l'ID de l'événement comme clé naturelle
    _, 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)
    }

    // Vérifier si l'insertion a réellement eu lieu (1 ligne) ou s'il s'agit d'un non-opération (0 lignes)
    // Une approche plus simple : utiliser RETURNING ou vérifier les lignes affectées
    // Si 0 ligne affectée, il s'agit d'un doublon -- passez-le
    ...
}

En pratique, de nombreuses équipes s’appuient sur les propres en-têtes de déduplication du courtier (comme le champ key de Kafka pour les sujets à journalisation compactée, ou l’en-tête message-id de RabbitMQ) et considèrent la déduplication au niveau de la base de données comme une solution de secours. Les deux approches sont des couches valides à appliquer.

Incluez l’id de l’événement outbox (un UUID) dans le message publié en tant que clé de déduplication. Les consommateurs peuvent alors l’utiliser indépendamment du mécanisme de déduplication qu’ils préfèrent.

Politique de nouvelle tentative et messages empoisonnés

La colonne attempts pilote la politique de nouvelle tentative. Le relais saute les lignes où attempts >= maxAttempts et traite ces lignes comme des lettres mortes (dead letters). Un processus distinct ou une alerte opérateur les gère.

Une vue simple de lettre morte :

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

Une bonne politique de nouvelle tentative en production :

  • Définissez maxAttempts à 5-10 en fonction du coût des nouvelles tentatives.
  • Envisagez une backoff exponentielle : incluez une colonne retry_after et sautez les lignes où retry_after > NOW().
  • Alerte sur COUNT(*) FROM outbox_dead_letters dépassant un seuil.
  • Fournissez un chemin de nouvelle tentative manuel : un point de terminaison administrateur ou un script qui réinitialise attempts = 0 et retry_after = NULL pour des lignes spécifiques.

Les messages empoisonnés – les lignes qui échouent systématiquement en raison d’un bug dans le consommateur ou d’une incompatibilité de schéma – ne doivent pas bloquer les messages sains. Puisque le relais traite un lot par intervalle et marque les échecs par une incrémentation d’essai plutôt que de les supprimer de la file d’attente, les lignes saines se déroulent normalement tandis que celles empoisonnées accumulent des essais jusqu’à ce qu’elles atteignent le seuil de lettre morte. Cette vue outbox_dead_letters est une version côté base de données du même modèle files d’attente de lettres mortes implémentées nativement par les courtiers – quarantaine après un seuil, alerte sur le volume, et décision délibérée requise avant le rejeu.

Tri des événements et partitionnement

La requête de sondage trie par created_at, ce qui donne un ordre de premier arrivé, premier servi (FIFO) au sein d’un lot. Pour la plupart des cas d’utilisation, cela suffit. Lorsque le tri par entité est strictement nécessaire – par exemple, s’assurer que order.updated n’est jamais publié avant order.created pour la même commande – vous avez besoin d’un tri par agrégat.

Ajoutez aggregate_id à la clause ORDER BY et utilisez-le comme clé de message lors de la publication vers un sujet partitionné comme Apache Kafka. Kafka route tous les messages avec la même clé vers la même partition, et les partitions sont consommées dans l’ordre. Cela vous offre des garanties de tri par agrégat sans tri global, ce qui nécessiterait une seule instance de relais.

ORDER BY aggregate_id, created_at

Pour les courtiers qui ne supportent pas le tri partitionné (comme les files d’attente AMQP de base), le relais à instance unique ou les vérifications de tri au niveau de l’application dans le consommateur sont les alternatives pratiques.

Réduction de la latence de sondage avec LISTEN/NOTIFY

Un intervalle de sondage d’une seconde signifie une latence moyenne des événements de 500 millisecondes. Pour la plupart des charges de travail, c’est acceptable. Pour les cas où vous avez besoin d’une latence quasi nulle, le mécanisme LISTEN/NOTIFY de PostgreSQL permet au relais de se réveiller immédiatement lorsqu’une nouvelle ligne outbox est insérée.

Ajoutez un déclencheur à la table 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();

Dans le relais, écoutez sur le canal et réveillez-vous sur les notifications tout en conservant le sondage périodique comme solution de secours :

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) // sondage de secours
    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)
            }
        }
    }
}

Le temporisateur de secours gère toute notification manquée lors d’un redémarrage du relais ou d’un hiccup réseau. Gardez l’intervalle de secours à quelques secondes plutôt qu’à des millisecondes – sa tâche est la récupération, pas la faible latence.

Observabilité : métriques, journaux et alertes

L’outbox est une infrastructure. Traitez-la comme telle et instrumentez-la en conséquence.

Métriques clés :

var (
    outboxPublished = prometheus.NewCounter(prometheus.CounterOpts{
        Name: "outbox_events_published_total",
        Help: "Total des événements outbox publiés avec succès.",
    })
    outboxFailed = prometheus.NewCounterVec(prometheus.CounterOpts{
        Name: "outbox_events_failed_total",
        Help: "Total des échecs de publication outbox par type d'événement.",
    }, []string{"event_type"})
    outboxPending = prometheus.NewGauge(prometheus.GaugeOpts{
        Name: "outbox_events_pending",
        Help: "Nombre actuel d'événements outbox non traités.",
    })
    outboxBatchDuration = prometheus.NewHistogram(prometheus.HistogramOpts{
        Name:    "outbox_batch_duration_seconds",
        Help:    "Durée de chaque lot de traitement outbox.",
        Buckets: prometheus.DefBuckets,
    })
)

Actualisation des jauges : exécutez une requête périodique pour maintenir outbox_events_pending précis :

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

Seuils d’alerte à envisager :

  • outbox_events_pending > 1000 pendant plus de deux minutes : le relais prend du retard ou est bloqué.
  • outbox_events_pending croissant de manière monotone : le courtier est indisponible ou le relais a planté.
  • Le nombre de lettres mortes non nul : un bug de schéma ou de consommateur nécessite une investigation.
  • outbox_batch_duration_seconds p95 > 5s : la base de données est lente ou la taille du lot est trop grande.

Champs de journal structurés : incluez event_id, event_type, aggregate_id et attempt dans chaque ligne de journal du relais. Ces champs vous permettent de corrélérer une publication échouée avec la ligne outbox spécifique et la trace du consommateur en aval.

Outbox vs. file d’attente directe vs. Saga

Le modèle outbox n’est pas l’outil approprié pour chaque problème de coordination. Voici la comparaison :

Approche Atomicité Complexité Quand l’utiliser
Publication directe Aucune Faible Acceptable de perdre occasionnellement des événements
Outbox transactionnelle Forte Moyenne Livraison fiable d’événements depuis un seul service
Modèle Saga Finalement Élevée Transactions multi-services qui s’étendent sur plusieurs bases de données
Commit en deux phases Forte Très élevée Rarement pratique ; évité dans la plupart des systèmes distribués

Le modèle outbox garantit qu’un service émet de manière fiable des événements qui reflètent ses propres changements d’état. Il ne coordonne pas les changements d’état entre plusieurs services – c’est ce que le Modèle Saga est conçu pour le faire. Le choix du courtier – qu’il s’agisse de RabbitMQ, SQS, ou Kafka – est indépendant du modèle outbox lui-même ; le relais publie vers n’importe quel courtier que votre système utilise.

Si vous construisez une saga, le modèle outbox reste utile : chaque participant à la saga écrit son changement d’état local et son événement saga dans une transaction unique en utilisant l’outbox, puis l’orchestrateur ou la chorégraphie de la saga lit ces événements de manière fiable.

CDC basé sur les journaux de transactions (WAL) comme alternative de relais

Au lieu de faire du sondage, vous pouvez suivre le journal de transactions (Write-Ahead Log - WAL) de PostgreSQL et lire les insertions outbox directement depuis le flux de réplication. Des outils comme Debezium font cela. Les avantages sont une latence plus faible et aucune pression de verrouillage sur la table outbox. Les inconvénients sont la complexité opérationnelle, un slot de réplication PostgreSQL dédié et un service externe à exécuter et à surveiller.

Pour la plupart des équipes, le relais par sondage décrit ci-dessus est le bon point de départ. Le suivi du WAL a du sens lorsque vous avez des taux d’insertion outbox élevés (des dizaines de milliers par seconde), une latence d’événement inférieure à 100 ms, ou que vous utilisez déjà Debezium pour d’autres besoins de capture de changements.

Intégration sqlc

Si vous utilisez sqlc pour du code Go de base de données typé en toute sécurité, les requêtes outbox s’intègrent naturellement :

-- 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 génère des fonctions typées en toute sécurité pour chaque requête, ce qui évite les erreurs d’interpolation de chaînes et maintient la logique de requête outbox au même endroit que le reste de votre couche d’accès aux bases de données.

Liste de vérification de production

Utilisez ceci avant de déployer une implémentation outbox :

Base de données

  • La table outbox a l’index partiel sur created_at WHERE processed_at IS NULL
  • La colonne attempts est présente avec une valeur par défaut de 0
  • Une vue ou une requête de lettre morte est définie
  • Les anciennes lignes traitées sont archivées ou supprimées périodiquement (un job de nettoyage nocturne suffit)

Relais

  • FOR UPDATE SKIP LOCKED est utilisé dans la requête de sondage
  • Le relais s’exécute à l’intérieur d’une transaction (début avant la requête, validation après toutes les mises à jour)
  • La taille du lot est bornée (50-200 lignes est typique)
  • Le relais respecte l’annulation du contexte pour l’arrêt progressif
  • Les publications échouées incrémente attempts au lieu de provoquer l’annulation du lot

Idempotence

  • Le message publié inclut l’id outbox en tant que clé de déduplication
  • Les consommateurs sont idempotents ou le courtier fournit la déduplication
  • Consultez Idempotence dans les Systèmes Distribués pour les modèles de déduplication

Observabilité

  • La jauge outbox_events_pending est surveillée et des alertes sont configurées
  • Le nombre de lettres mortes est surveillé par alerte
  • La durée du lot du relais est suivie
  • Les journaux structurés incluent event_id, event_type et aggregate_id

Opérations

  • Un chemin de nouvelle tentative manuel existe pour les lignes de lettre morte
  • Le comportement de redémarrage du relais est testé (re-publie-t-il correctement ?)
  • Le comportement en cas d’indisponibilité du courtier est testé (l’outbox grandit et se vide-t-elle correctement ?)

Réflexions finales

Le problème du double écriture est facile à rejeter comme un cas边缘 (edge case) jusqu’à ce qu’il cause un incident. Le modèle de la boîte aux lettres transactionnelle le résout avec des outils que vous avez déjà : une transaction PostgreSQL, une goroutine en arrière-plan et une table supplémentaire. Le relais est simple à construire, simple à exploiter et simple à comprendre.

Le coût est que les consommateurs doivent être conçus pour une livraison au moins une fois (at-least-once). C’est un compromis raisonnable. Une livraison exactement une fois (exactly-once) entre une base de données et un courtier sans transactions distribuées n’est pas réalisable en pratique – et prétendre le contraire conduit à des systèmes qui perdent ou traitent en double silencieusement des événements dans des conditions d’échec.

Écrivez l’événement avec les données. Relayez-le de manière fiable. Rendez les consommateurs idempotents. C’est tout le modèle.

Cet article fait partie du cluster Architecture d’Application en Production.

Sources

S'abonner

Recevez de nouveaux articles sur les systèmes, l'infrastructure et l'ingénierie IA.