Scénarios et exemples de files d'attente Spanner

Ce document fournit des modèles d'architecture et des exemples de code pour les scénarios de messagerie courants utilisant des files d'attente Spanner. Vous pouvez utiliser ces modèles pour déclencher des tâches asynchrones après la validation des transactions, planifier des tâches différées ou récurrentes, gérer de grandes charges utiles de messages avec un stockage hors bande, coordonner des workflows multi-événements, et créer des points de contrôle ou prolonger des baux pour les jobs de longue durée en arrière-plan.

Traitement "exactement une fois" et accusé de réception "au maximum une fois"

Les différentes considérations et solutions pour le traitement de type "exactement une fois" et l'accusé de réception de type "au maximum une fois" sont décrites plus en détail sur la page Traitement de type "exactement une fois" et accusé de réception de type "au maximum une fois".

Effectuer des tâches après la validation d'une transaction

Pour effectuer des tâches après la validation d'une transaction, envoyez un message à la file d'attente au cours de la même transaction.

Par exemple, l'inscription d'un nouvel utilisateur déclenche l'envoi d'un e-mail de bienvenue :

GoogleSQL

-- Inside your application transaction:
-- 1. Insert into Users table
INSERT INTO Users (UserId, UserName) VALUES (124, 'New User');

-- 2. Send message to queue to trigger email
INSERT INTO UserTasks (UserId, MessageId, Payload)
VALUES (
  124,
  'welcome-email-id',
  b'{"type": "welcome", "email": "user@example.com"}'
);

PostgreSQL

-- Inside your application transaction:
-- 1. Insert into users table
INSERT INTO users (userid, username) VALUES (124, 'New User');

-- 2. Send message to queue to trigger email
INSERT INTO usertasks (userid, messageid, payload)
VALUES (
  124,
  'welcome-email-id',
  CAST('{"type": "welcome", "email": "user@example.com"}' AS bytea)
);

Une fois la transaction validée, le récepteur de UserTasks diffuse le message, envoie l'e-mail et accuse réception du message :

GoogleSQL

-- 1. In the receiver process, stream messages from the queue
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken
FROM RECEIVE_UserTasks(max_duration=>'20m');

-- 2. After sending the welcome email, acknowledge the message
DELETE FROM UserTasks
WHERE UserId = 124 AND MessageId = 'welcome-email-id';

PostgreSQL

-- 1. In the receiver process, stream messages from the queue
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token
FROM spanner.receive_usertasks(NULL, NULL, '20m');

-- 2. After sending the welcome email, acknowledge the message
DELETE FROM usertasks
WHERE userid = 124 AND messageid = 'welcome-email-id';

Gérer les tâches de longue durée

Si vous avez une tâche qui peut prendre plus de temps que la durée de bail par défaut (plus de 10 secondes), appelez SELECT * FROM RENEWLEASE_QUEUE_NAME() régulièrement.

Par exemple, pour générer un rapport :

  1. Le destinataire reçoit un message de RECEIVE_ReportQueue().
  2. Lancez la génération du rapport.
  3. Appelez SELECT * FROM RENEWLEASE_ReportQueue([leaseToken]) toutes les cinq secondes dans un thread ou une routine distincts.
  4. Une fois l'opération terminée, accusez réception du message et stockez le rapport.

Si vous avez une tâche de longue durée qui nécessite un traitement "au plus une fois" ou une longue durée de bail, procédez comme suit :

  1. Accusez réception (DELETE ou ACK) du message actuel de la file d'attente à son arrivée. Dans la même transaction, remettez en file d'attente un nouveau message avec un code temporel de remise futur, au-delà du temps nécessaire au traitement.
  2. Poursuivez le traitement et confirmez le nouveau message mis en file d'attente une fois terminé.

Cette approche présente l'avantage de ne pas nécessiter d'étendre continuellement le bail, et le message n'est pas renvoyé tant que l'heure future n'est pas arrivée (ce qui couvre les plantages). Si l'accusé de réception initial réussit, le traitement de type "au plus une fois" est effectué.

Créer des points de contrôle pour les tâches de longue durée

Les files d'attente Spanner peuvent gérer des tâches qui durent de quelques minutes à plusieurs heures, et pas seulement des tâches rapides. Pour ces tâches de longue durée, utilisez l'approche suivante :

  1. Stocker les métadonnées en externe : utilisez le stockage hors bande pour conserver les détails et l'état de la tâche.
  2. Définissez régulièrement des points de contrôle : pour récupérer après un plantage sans perdre trop de progression, la tâche doit enregistrer régulièrement son état.
  3. Utilisez le modèle de point de contrôle recommandé : la meilleure façon de créer un point de contrôle consiste à accuser réception de manière atomique (ACK) du message de file d'attente actuel et à envoyer un nouveau message dont la diffusion est planifiée. Ce nouveau message contient ou pointe vers l'état mis à jour, ce qui empêche une nouvelle remise immédiate à un autre employé.

Ce modèle réduit le travail en double, même si la création de points de contrôle complets n'est pas possible. Toutefois, dans ce scénario, la tâche redémarre depuis le début après un plantage.

Planifier une tâche pour une heure spécifique dans le futur

Pour planifier une tâche à une heure spécifique, définissez la colonne DeliverTime lorsque vous insérez le message.

Par exemple, un rappel de l'expiration d'un essai :

GoogleSQL

-- 1. Insert into Users table
INSERT INTO Users (UserId, UserName) VALUES (125, 'Trial User');

-- 2. Send message to queue with a future delivery time
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (125, 'trial-expire-reminder', b'{"type": "reminder"}', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 29 DAY));

PostgreSQL

-- 1. Insert into users table
INSERT INTO users (userid, username) VALUES (125, 'Trial User');

-- 2. Send message to queue with a future delivery time
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (125, 'trial-expire-reminder', CAST('{"type": "reminder"}' AS bytea), CURRENT_TIMESTAMP + INTERVAL '29 DAY');

Gérer les charges utiles de messages volumineuses

Si la charge utile de votre message est volumineuse, utilisez le modèle de stockage hors bande. Stockez la charge utile volumineuse dans une table distincte et ajoutez-y une référence dans le message de la file d'attente.

Par exemple, le traitement d'images :

GoogleSQL

-- Schema
CREATE TABLE ImageUploads (
  UserId    INT64 NOT NULL,
  ImageId   STRING(36) NOT NULL,
  ImageData BYTES(MAX),
  Status    STRING(MAX) -- PENDING, PROCESSING, DONE
) PRIMARY KEY (UserId, ImageId),
  INTERLEAVE IN PARENT Users;

CREATE QUEUE ImageProcessingQueue (
  UserId    INT64 NOT NULL,
  ImageId   STRING(36) NOT NULL,
  Payload   BYTES(1) NOT NULL -- Payload can be minimal
) PRIMARY KEY (UserId, ImageId),
  INTERLEAVE IN PARENT ImageUploads ON DELETE CASCADE;

-- Application Logic
-- 1. Upload image, insert into ImageUploads with Status 'PENDING'
-- 2. Send message to ImageProcessingQueue
INSERT INTO ImageProcessingQueue (UserId, ImageId, Payload) VALUES (123, 'image-uuid-1', b'');

-- Receiver for ImageProcessingQueue:
-- 1. Receives message (UserId, ImageId).
-- 2. Reads ImageData from ImageUploads.
-- 3. Processes image.
-- 4. Updates ImageUploads Status to 'DONE'.
-- 5. ACKs the queue message.

PostgreSQL

-- Schema
CREATE TABLE imageuploads (
  userid    bigint NOT NULL,
  imageid   varchar(36) NOT NULL,
  imagedata bytea,
  status    varchar, -- PENDING, PROCESSING, DONE
  PRIMARY KEY (userid, imageid)
) INTERLEAVE IN PARENT users;

CREATE QUEUE imageprocessingqueue (
  userid    bigint NOT NULL,
  imageid   varchar(36) NOT NULL,
  payload   bytea NOT NULL, -- Payload can be minimal
  PRIMARY KEY (userid, imageid)
) INTERLEAVE IN PARENT imageuploads ON DELETE CASCADE;

-- Application Logic
-- 1. Upload image, insert into imageuploads with status 'PENDING'
-- 2. Send message to imageprocessingqueue
INSERT INTO imageprocessingqueue (userid, imageid, payload) VALUES (123, 'image-uuid-1', CAST('' AS bytea));

-- Receiver for imageprocessingqueue:
-- 1. Receives message (userid, imageid).
-- 2. Reads imagedata from imageuploads.
-- 3. Processes image.
-- 4. Updates imageuploads status to 'DONE'.
-- 5. ACKs the queue message.

Attendre plusieurs événements avant de continuer

Pour attendre plusieurs événements avant de continuer (comme une opération de jointure), utilisez un tableau pour suivre l'état et une file d'attente pour déclencher les vérifications.

Par exemple, le traitement des commandes nécessitant un inventaire et un paiement :

  1. Créez une table Orders avec InventoryStatus et PaymentStatus.
  2. Une fois l'inventaire confirmé, mettez à jour Orders et envoyez un message à OrderCheckQueue.
  3. Une fois le paiement confirmé, mettez à jour Orders et envoyez un message à OrderCheckQueue.
  4. Le récepteur de OrderCheckQueue vérifie la table Orders. Si les deux états sont confirmés, l'expédition est effectuée et le message est accusé de réception. Si ce n'est pas le cas, il peut être mis en file d'attente pour une vérification ultérieure ou exécuter une autre logique.

Effectuer une action périodiquement

Pour effectuer une action périodiquement, utilisez le modèle de planification périodique. Le destinataire accuse réception du message et en envoie un nouveau planifié pour le prochain intervalle.

Par exemple, l'agrégation des données horaires :

GoogleSQL

-- Inside your application transaction:
-- 1. Acknowledge current message
DELETE FROM AggregationQueue
WHERE TaskType = 'hourly-aggregator' AND MessageId = 'current-uuid'
ASSERT_ROWS_MODIFIED 1;

-- 2. Schedule next run 1 hour in the future
INSERT INTO AggregationQueue (TaskType, MessageId, Payload, DeliverTime)
VALUES ('hourly-aggregator', 'next-uuid', b'', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR));

PostgreSQL

-- Inside your application transaction:
-- 1. Acknowledge current message
DELETE FROM aggregationqueue
WHERE tasktype = 'hourly-aggregator' AND messageid = 'current-uuid'
ASSERT_ROWS_MODIFIED 1;

-- 2. Schedule next run 1 hour in the future
INSERT INTO aggregationqueue (tasktype, messageid, payload, deliver_time)
VALUES ('hourly-aggregator', 'next-uuid', CAST('' AS bytea), CURRENT_TIMESTAMP + INTERVAL '1 HOUR');

Vous pouvez également utiliser les mutations Ack et Send de la bibliothèque cliente. Ces exemples supposent que vous disposez d'un objet Message encapsulant la clé et la charge utile :

Java

// Receiver logic for AggregationQueue
public void process(DatabaseClient dbClient, Message msg) {
  // ... do aggregation ...

  // ACK current message and schedule next run (1 hour from now)
  Instant nextRun = Instant.now().plus(Duration.ofHours(1));
  Mutation ackMutation =
      Mutation.newAckBuilder("AggregationQueue")
          .setKey(msg.getKey()) // Ack
          .build();
  Mutation sendMutation =
      Mutation.newSendBuilder("AggregationQueue")
          .setKey(Key.of("hourly-aggregator", "next-uuid"))
          .setPayload(Value.bytes(ByteArray.copyFrom("")))
          .setDeliveryTime(nextRun) // Schedule next
          .build();
  dbClient.write(Arrays.asList(ackMutation, sendMutation));
}

Go

// Receiver logic for AggregationQueue
func process(msg) {
    // ... do aggregation ...

    // ACK current message and schedule next run
    nextRun := time.Now().Add(1 * time.Hour)
    _, err := client.Apply(ctx, []*spanner.Mutation{
        spanner.Ack("AggregationQueue", msg.Key), // Ack
        spanner.Send("AggregationQueue",
            spanner.Key{"hourly-aggregator", "next-uuid"},
            []byte(""),
            spanner.WithDeliveryTime(nextRun), // Schedule next
        ),
    })
    // ... handle err ...
}

Python

# Receiver logic for AggregationQueue
def process(database: spanner.Database, msg: Message):
  # ... do aggregation ...
  # ACK current message and schedule next run (1 hour from now)
  next_run = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(
      hours=1
  )
  with database.batch() as batch:
    batch.ack(
        queue="AggregationQueue",
        key=msg.key,  # Ack
    )
    batch.send(
        queue="AggregationQueue",
        key=("hourly-aggregator", "next-uuid"),
        payload=b"",
        deliver_time=next_run,  # Schedule next
    )

Node.js

/**
 * Receiver logic for AggregationQueue
 * @param {import('@google-cloud/spanner').Database} database
 * @param { { key: Array<string|number>, payload: Buffer } } msg
 */
async function process(database, msg) {
  // ... do aggregation ...
  // ACK current message and schedule next run (1 hour from now)
  const nextRun = new Date(Date.now() + 60 * 60 * 1000);
  await database.runTransactionAsync(async (transaction) => {
    // Ack current message
    transaction.queueAck('AggregationQueue', msg.key);
    // Schedule next run
    transaction.queueSend(
      'AggregationQueue',
      ['hourly-aggregator', 'next-uuid'],
      {
        payload: Buffer.from(''),
        deliverTime: nextRun,
      }
    );
    await transaction.commit();
  });
}

Regrouper les messages par lot à l'aide du traitement par lot temporel

Les files d'attente Spanner distribuent les messages avec une latence minimale. Toutefois, lorsque de grands volumes de messages arrivent en continu de clients indépendants, le traitement individuel de chaque message peut entraîner des frais de transaction élevés. Si vous essayez d'interroger ou d'analyser manuellement la table de file d'attente pour regrouper les messages, vous risquez de provoquer des conflits de verrouillage de plage, d'augmenter le taux d'abandon et d'entraîner des coûts de lecture supplémentaires.

Pour obtenir un traitement par lot à haut débit sans conflit, appliquez le modèle de traitement par lot temporel. Les expéditeurs alignent le DeliverTime des messages sur une période discrète dans le futur (par exemple, en arrondissant à la limite de 10 secondes la plus proche). Étant donné que les expéditeurs indépendants calculent un code temporel futur identique, Spanner a tendance à regrouper les messages d'une même division et à les distribuer dans un seul lot si max_batch_size dans la fonction table RECEIVE_QUEUE_NAME() le permet, ou dans plusieurs lots si max_batch_size est inférieur au nombre de messages à distribuer.

Le code temporel de la diffusion peut être calculé à l'aide de la formule suivante :

$$ \text{DeliverTime} = \text{RoundDown}(\text{CurrentTime}, \text{FixedDelay}) + \text{FixedDelay} $$

Par exemple, avec une fenêtre de 10 secondes, les messages mis en file d'attente entre 09:05:00 et 09:05:09.999 reçoivent tous un DeliverTime de 09:05:10 :

GoogleSQL

-- Calculate delivery time rounded to the next 10-second interval.
-- Use DIV(..., 10) * 10 to perform the RoundDown in SQL:
INSERT INTO OrderProcessingQueue (OrderId, DeliverTime, Payload)
VALUES (
  'order-101',
  TIMESTAMP_SECONDS(DIV(UNIX_SECONDS(CURRENT_TIMESTAMP()), 10) * 10 + 10),
  b'{"item": "book", "qty": 1}'
);

PostgreSQL

-- Calculate delivery time rounded to the next 10-second interval:
INSERT INTO orderprocessingqueue (orderid, deliver_time, payload)
VALUES (
  'order-101',
  to_timestamp((floor(extract(epoch from CURRENT_TIMESTAMP) / 10) * 10) + 10),
  CAST('{"item": "book", "qty": 1}' AS bytea)
);

Les destinataires récupèrent ensuite ces messages synchronisés par lots à l'aide de max_batch_size :

GoogleSQL

SELECT OrderId, Payload, DeliverTime, SpannerLeaseToken
FROM RECEIVE_OrderProcessingQueue(max_duration=>'20m', max_batch_size=>50);

PostgreSQL

SELECT orderid, payload, deliver_time, spanner_lease_token
FROM spanner.receive_orderprocessingqueue(
    max_batch_size=>50, priority=>NULL, max_duration=>'20m');

Si vous envoyez des millions de messages en même temps, l'alignement de tous les messages sur la même seconde exacte peut entraîner des pics de traitement soudains. Pour répartir le travail de manière uniforme tout en regroupant les messages par entité, ajoutez un décalage au calcul en fonction d'un identifiant unique (tel que TIMESTAMP_ADD(deliver_time, INTERVAL MOD(ABS(FARM_FINGERPRINT(CAST(UserId AS STRING))), 10) SECOND)).

Dissocier la modification des données du traitement (modèle d'indicateur de données modifiées)

Dans les applications transactionnelles à haut débit, l'exécution de recomputations complexes, d'indexation de recherche ou d'invalidation de cache directement dans les transactions destinées aux utilisateurs peut ralentir l'expérience utilisateur en augmentant la latence et en provoquant des conflits de verrouillage sur les lignes partagées.

Le modèle d'indicateur de modification découple les modifications de données du traitement asynchrone. Lorsqu'une transaction modifie une table, elle écrit un message "bit sale" léger dans une file d'attente au sein de la même transaction. Un nœud de calcul en arrière-plan consomme ensuite le message et effectue le traitement coûteux de manière asynchrone.

Par exemple, lorsqu'un client modifie son profil ou ses paramètres :

GoogleSQL

-- Inside user profile update transaction:
-- 1. Update the primary entity table
UPDATE UserProfiles
SET FullName = 'Jane Doe', UpdatedAt = CURRENT_TIMESTAMP()
WHERE UserId = 456;

-- 2. Send lightweight dirty flag message to the queue.
-- It is recommended to interleave the queue in the primary UserProfiles
-- table for better transaction performance.
INSERT INTO UserDirtyQueue (UserId, TaskType, CommitTimestamp, Payload)
VALUES (456, 'reindex-user-profile', CURRENT_TIMESTAMP(), b'');

PostgreSQL

-- Inside user profile update transaction:
-- 1. Update the primary entity table
UPDATE userprofiles
SET fullname = 'Jane Doe', updatedat = CURRENT_TIMESTAMP
WHERE userid = 456;

-- 2. Send lightweight dirty flag message to the queue
INSERT INTO userdirtyqueue (userid, tasktype, committimestamp, payload)
VALUES (456, 'reindex-user-profile', CURRENT_TIMESTAMP, CAST('' AS bytea));

Le récepteur d'arrière-plan pour UserDirtyQueue reçoit le UserId, lit la ligne de profil actualisée en dehors du chemin critique de l'utilisateur et recalcule l'index de recherche ou met à jour les caches externes. Le traitement par lot peut également être appliqué à ce modèle si plusieurs mises à jour concernant le même utilisateur sont envoyées à la file d'attente dans un court laps de temps. Dans ce cas, le TVF du récepteur peut spécifier un max_batch_size supérieur à 1 pour recevoir plusieurs messages du même lot.

Surveiller l'état des workers et détecter les délais d'attente

Vous pouvez utiliser des messages de file d'attente planifiés pour créer un système de surveillance de l'état et de signal de présence tolérant aux pannes pour les flottes de nœuds de calcul ou les instances de microservices.

Pour implémenter la vérification de l'état :

  1. Enregistrer le nœud de calcul au démarrage : lorsqu'un nœud de calcul s'initialise, il insère un message de signal de présence dans une file d'attente de vérification de l'état avec un DeliverTime futur défini sur son délai d'échec (par exemple, 60 secondes).
  2. Envoyer des signaux de présence périodiques : tant qu'il est en bon état, le nœud de calcul actualise périodiquement son signal de présence (par exemple, toutes les 10 secondes) en avançant le DeliverTime de 60 secondes.
  3. Détecter les échecs : si le nœud de calcul plante ou perd la connectivité réseau, les actualisations du signal de présence s'arrêtent. Au bout de 60 secondes, le code temporel de distribution devient valide (DeliverTime <= CURRENT_TIMESTAMP()) et Spanner distribue le message à un destinataire d'alerte, qui lance le basculement ou la réattribution de la tâche.

Important : Les files d'attente Spanner ne sont pas compatibles avec les instructions UPDATE LMD. Par conséquent, pour actualiser le code temporel du signal de présence, vous devez supprimer le message existant et insérer un message de remplacement avec le nouveau DeliverTime dans une seule transaction, ou appliquer les mutations Ack et Send de la bibliothèque cliente. Assurez-vous que la clé primaire de la file d'attente est WorkerId seul (plutôt que (WorkerId, MessageId)) afin qu'il n'y ait qu'un seul message de signal de présence par nœud de calcul à un moment donné.

GoogleSQL

-- Inside the worker heartbeat transaction (executed every 10 seconds):
-- 1. Acknowledge the existing heartbeat message
DELETE FROM WorkerHealthQueue
WHERE WorkerId = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

-- 2. Send replacement heartbeat with refreshed 60-second deadline
INSERT INTO WorkerHealthQueue (WorkerId, DeliverTime, Payload)
VALUES (
  'worker-node-42',
  TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 60 SECOND),
  b'{"status": "healthy", "active_jobs": 3}'
);

PostgreSQL

-- Inside the worker heartbeat transaction (executed every 10 seconds):
-- 1. Acknowledge the existing heartbeat message
DELETE FROM workerhealthqueue
WHERE workerid = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

-- 2. Send replacement heartbeat with refreshed 60-second deadline
INSERT INTO workerhealthqueue (workerid, deliver_time, payload)
VALUES (
  'worker-node-42',
  CURRENT_TIMESTAMP + INTERVAL '60 SECOND',
  CAST('{"status": "healthy", "active_jobs": 3}' AS bytea)
);

Lorsque vous utilisez des mutations de bibliothèque cliente, assurez-vous que Ack précède Send dans le slice de mutation, comme dans cet exemple Go :

// Worker heartbeat loop
func sendHeartbeat(ctx context.Context, client *spanner.Client, workerID string) error {
    newDeadline := time.Now().Add(60 * time.Second)
    _, err := client.Apply(ctx, []*spanner.Mutation{
        spanner.Ack("WorkerHealthQueue", spanner.Key{workerID}),
        spanner.Send(
            "WorkerHealthQueue",
            spanner.Key{workerID},
            []byte(`{"status":"healthy"}`),
            spanner.WithDeliveryTime(newDeadline),
        ),
    })
    return err
}

Lorsqu'un nœud de calcul s'arrête de manière ordonnée, il supprime explicitement son message de signal de présence afin qu'aucune fausse alerte ne soit déclenchée :

DELETE FROM WorkerHealthQueue WHERE WorkerId = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

Gérer plusieurs types de tâches dans une même file d'attente

Les instances Spanner sont soumises à des limites concernant le nombre total de files d'attente. La création d'une file d'attente distincte pour chaque petite opération asynchrone peut rapidement atteindre cette limite et nécessite de gérer de nombreuses requêtes de récepteur simultanées.

Pour consolider les opérations, combinez différents types de tâches dans une seule file d'attente, appelée file d'attente polymorphe. Deux stratégies permettent de créer une file d'attente polymorphe.

Stratégie 1 : Inclure une colonne de type dans la clé primaire

GoogleSQL

CREATE QUEUE ApplicationTasks (
  TaskType   STRING(50) NOT NULL,
  TaskId     STRING(36) NOT NULL,
  Payload    BYTES(MAX) NOT NULL,
) PRIMARY KEY (TaskType, TaskId);

-- Enqueue an email task
INSERT INTO ApplicationTasks (TaskType, TaskId, Payload)
VALUES ('SEND_EMAIL', 'task-uuid-1', b'{"to": "user@example.com", "template": "welcome"}');

-- Enqueue an image thumbnail task
INSERT INTO ApplicationTasks (TaskType, TaskId, Payload)
VALUES ('GENERATE_THUMBNAIL', 'task-uuid-2', b'{"image_id": "img-789", "size": "small"}');

PostgreSQL

CREATE QUEUE applicationtasks (
  tasktype   varchar(50) NOT NULL,
  taskid     varchar(36) NOT NULL,
  payload    bytea NOT NULL,
  PRIMARY KEY (tasktype, taskid)
);

-- Enqueue an email task
INSERT INTO applicationtasks (tasktype, taskid, payload)
VALUES ('SEND_EMAIL', 'task-uuid-1', CAST('{"to": "user@example.com", "template": "welcome"}' AS bytea));

-- Enqueue an image thumbnail task
INSERT INTO applicationtasks (tasktype, taskid, payload)
VALUES ('GENERATE_THUMBNAIL', 'task-uuid-2', CAST('{"image_id": "img-789", "size": "small"}' AS bytea));

Le récepteur inspecte TaskType et distribue la charge utile au gestionnaire correspondant.

Stratégie 2 : Structure de charge utile polymorphe

Vous pouvez également utiliser une charge utile JSON contenant un champ d'action ou de type discriminatoire :

{
  "action": "SYNC_INVENTORY",
  "data": { "item_id": 987, "delta": -1 }
}

Implémenter des délais de nouvelle tentative personnalisés

Les files d'attente Spanner relancent automatiquement les messages ayant échoué ou non reconnus avec un intervalle exponentiel entre les tentatives intégré. Toutefois, dans les scénarios où un message échoue pour une raison connue avec une durée connue ou une limite de débit externe (par exemple, une réponse HTTP 429 spécifiant un en-tête Retry-After), s'appuyer sur l'intervalle entre les tentatives automatique peut entraîner des tentatives prématurées qui gaspillent les ressources du processeur.

Pour implémenter un délai de nouvelle tentative personnalisé :

  1. Détectez l'échec temporaire spécifique dans votre processeur de messages.
  2. Confirmez le message actuel pour satisfaire la tentative de remise actuelle.
  3. Dans la même transaction, envoyez un message de remplacement avec un DeliverTime explicite défini sur le délai de nouvelle tentative choisi (dans l'exemple suivant, 5 minutes plus tard).

GoogleSQL

-- Inside failure-handling transaction:
-- 1. Acknowledge the failed message
DELETE FROM OutboundNotificationQueue
WHERE NotificationId = 'notif-555' ASSERT_ROWS_MODIFIED 1;

-- 2. Reschedule delivery 5 minutes in the future
INSERT INTO OutboundNotificationQueue (NotificationId, DeliverTime, Payload)
VALUES (
  'notif-555',
  TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 5 MINUTE),
  b'{"recipient": "user@example.com", "retry_count": 2}'
);

PostgreSQL

-- Inside failure-handling transaction:
-- 1. Acknowledge the failed message
DELETE FROM outboundnotificationqueue
WHERE notificationid = 'notif-555' ASSERT_ROWS_MODIFIED 1;

-- 2. Reschedule delivery 5 minutes in the future
INSERT INTO outboundnotificationqueue (notificationid, deliver_time, payload)
VALUES (
  'notif-555',
  CURRENT_TIMESTAMP + INTERVAL '5 MINUTE',
  CAST('{"recipient": "user@example.com", "retry_count": 2}' AS bytea)
);

Cette approche permet à votre application de gérer précisément les plannings de backoff et d'éviter de saturer les API externes pendant les périodes de récupération en aval.

Étapes suivantes