Utiliser les files d'attente Spanner

Ce document explique comment utiliser les files d'attente Spanner. Il explique comment créer une file d'attente, envoyer et recevoir des messages, prolonger les baux de messages et accuser réception des messages. Il inclut également les bonnes pratiques, des informations sur la surveillance des files d'attente et des conseils de dépannage.

Créer une file d'attente

Pour créer une file d'attente, utilisez l'instruction CREATE QUEUE.

GoogleSQL

-- Example table for interleaving
CREATE TABLE Users (
  UserId   INT64 NOT NULL,
  UserName STRING(MAX)
) PRIMARY KEY (UserId);

-- Queue for processing user-related tasks
-- NOTE: Queues do not have to be interleaved, but for locality and
-- pre-warming purposes, it's recommended if you are inserting into a table
-- and a queue simultaneously.
CREATE QUEUE UserTasks (
  UserId     INT64 NOT NULL,
  MessageId  STRING(36) NOT NULL, -- UUID recommended
  Payload    BYTES(MAX) NOT NULL  -- Also: Proto, JSON, String are possible.
) PRIMARY KEY (UserId, MessageId),
INTERLEAVE IN PARENT Users ON DELETE CASCADE;

PostgreSQL

-- Example table for interleaving
CREATE TABLE users (
  userid   bigint NOT NULL,
  username varchar,
  PRIMARY KEY (userid)
);

-- Queue for processing user-related tasks
-- NOTE: Queues do not have to be interleaved, but for locality and
-- pre-warming purposes, it's recommended if you are inserting into a table
-- and a queue simultaneously.
CREATE QUEUE usertasks (
  userid     bigint NOT NULL,
  messageid  varchar(36) NOT NULL, -- UUID recommended
  payload    bytea NOT NULL, -- Also: text, varchar, jsonb are possible.
  PRIMARY KEY (userid, messageid)
) INTERLEAVE IN PARENT users ON DELETE CASCADE;

La seule colonne qui n'a pas besoin d'être créée explicitement dans l'instruction CREATE QUEUE s'appelle DeliverTime dans GoogleSQL et deliver_time dans PostgreSQL. Elles sont créées automatiquement par Spanner.

Les files d'attente sont compatibles avec les règles de durée de vie (TTL), qui peuvent vous aider à gérer le backlog de messages anciens et non confirmés.

Envoyer un message

Pour envoyer un message à une file d'attente, utilisez l'instruction LMD INSERT :

GoogleSQL

-- Send message immediately
INSERT INTO Users (UserId) VALUES (123);
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (123, 'some-unique-id-1', b'Your task payload here', CURRENT_TIMESTAMP());

-- Schedule a message delivery
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (123, 'some-unique-id-2', b'Scheduled task', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR));

PostgreSQL

-- Send message immediately
INSERT INTO users (userid) VALUES (123);
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (123, 'some-unique-id-1', CAST('Your task payload here' AS bytea), CURRENT_TIMESTAMP);

-- Schedule a message delivery
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (123, 'some-unique-id-2', CAST('Scheduled task' AS bytea), CURRENT_TIMESTAMP + INTERVAL '1 HOUR');

Vous pouvez également utiliser les mutations de la bibliothèque cliente Send pour insérer un message :

Go

m := spanner.Send("UserTasks", spanner.Key{int64(123), "some-unique-id-1"}, []byte("Your task payload here"), spanner.WithDeliveryTime(futureTime))
_, err := client.Apply(ctx, []*spanner.Mutation{m})

Java

dbClient.write(
        Collections.singletonList(
            Mutation.newSendBuilder("UserTasks")
                .setKey(Key.of(123L, "some-unique-id-1"))
                .setPayload(Value.bytes(ByteArray.copyFrom("message3")))
                .setDeliveryTime(futureTime)
                .build()));

Recevoir des messages

Utilisez la fonction de valeur de table (TVF) RECEIVE_QUEUE_NAME() avec ExecuteStreamingSQL pour recevoir des messages. Il s'agit d'un appel de longue durée. Vous devez exécuter l'un de ces appels par nœud de calcul et par file d'attente, de manière répétée. Exécutez la requête à l'aide d'une lecture forte, car Spanner rejette les lectures non actualisées.

GoogleSQL

-- SQL query to stream messages
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken
FROM RECEIVE_UserTasks(max_duration=>'20m');

PostgreSQL

-- SQL query to stream messages
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token
FROM spanner.receive_usertasks(NULL, NULL, '20m');

Votre code client doit parcourir les résultats à l'aide d'une requête de streaming. Chaque ligne renvoyée est un message.

Recevoir des messages par lot

Pour augmenter le débit en traitant plusieurs messages à la fois, vous pouvez recevoir des messages par lots en spécifiant l'argument max_batch_size :

GoogleSQL

-- SQL query to stream messages
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken,
  SpannerLastBatchMessage -- Special boolean column returns TRUE if the
                          -- last message is in a batch.
FROM RECEIVE_UserTasks(max_duration=>'20m', max_batch_size=>20);

PostgreSQL

-- SQL query to stream messages
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token,
  spanner_last_batch_message -- Special boolean column returns TRUE if the
                             -- last message is in a batch.
FROM spanner.receive_usertasks(20, NULL, '20m');

Vous pouvez également utiliser la bibliothèque cliente. Cet exemple Go montre comment diffuser des messages à partir d'une file d'attente, vérifier l'expiration du bail et accuser réception des messages de manière asynchrone :

// import "cloud.google.com/go/spanner"
// import "google.golang.org/api/iterator"

stmt := spanner.Statement{SQL: "SELECT * FROM RECEIVE_UserTasks(max_duration=>'20m')"}
iter := client.Single().Query(ctx, stmt)
defer iter.Stop()

for {
    row, err := iter.Next()
    if err == iterator.Done {
        break // Or potentially restart the query
    }
    if err != nil {
        // Handle error
        return err
    }

    var userId int64
    var messageId string
    var payload []byte
    var deliverTime time.Time
    var leaseExpiration time.Time
    var leaseToken string
    var lastBatchMessage bool

    if err := row.Columns(&userId, &messageId, &payload, &deliverTime, &leaseExpiration, &leaseToken, &lastBatchMessage); err != nil {
        // Handle column parsing error
        return err
    }

    if time.Now().After(leaseExpiration) {
        log.Printf("Lease expired for message %s, skipping", messageId)
        continue
    }

    // Process and acknowledge the message asynchronously so that we don't
    // block receiving subsequent messages.
    go func(userId int64, messageId string, payload []byte) {
        // Process the message (payload)
        // ... potentially long-running work ...
        // Need to extend the lease if processing is long

        // Acknowledge the message upon success
        _, err := client.Apply(ctx, []*spanner.Mutation{
            spanner.Ack("UserTasks", spanner.Key{userId, messageId}),
        })
        if err != nil {
            // Handle ack error
        }
    }(userId, messageId, payload)
}

Prolonger la durée de vie du message

Si le traitement d'un message prend plus de temps que le bail initial (10 secondes), utilisez la syntaxe suivante :

GoogleSQL

-- Extends the lease and returns new expiration time lease_token.
-- The lease will be extended by 10s (not configurable).
SELECT * FROM RENEWLEASE_UserTasks(lease_tokens => ['token1', ..., 'tokenN'])

-- Returns rows of tokens and whether they were successfully extended
SpannerOldLeaseToken   SpannerNewLeaseToken  SpannerLeaseExpirationTimestamp
        <old_token1>           <new_token1>  "2025-09-27T12:10:00.0Z"
            'token2'                 <NULL>  "2025-09-27T12:09:51.0Z"
...
            'tokenN'           'tokenN_new'  "2025-09-27T12:10:01.0Z"

PostgreSQL

-- Extends the lease and returns new expiration time lease_token.
-- The lease will be extended by 10s (not configurable).
SELECT * FROM spanner.renewlease_usertasks(lease_tokens => ARRAY['token1', ..., 'tokenN'])

-- Returns rows of tokens and whether they were successfully extended
spanner_old_lease_token   spanner_new_lease_token  spanner_lease_expiration_timestamp
           <old_token1>              <new_token1>  "2025-09-27T12:10:00.0Z"
               'token2'                    <NULL>  "2025-09-27T12:09:51.0Z"
...
               'tokenN'              'tokenN_new'  "2025-09-27T12:10:01.0Z"

Les jetons de bail sont renvoyés selon la logique suivante :

  1. Les jetons de bail non analysables ne renvoient pas de ligne.
  2. Les jetons de bail déjà expirés ne renvoient pas de ligne.
  3. Les jetons de bail non renouvelables renvoient une ligne avec une valeur SpannerNewLeaseToken définie sur NULL. Cela peut se produire si le message a déjà été confirmé, mais que le jeton de bail n'a pas expiré.

Vous pouvez également utiliser la bibliothèque cliente. Cet exemple Go montre comment étendre un bail de message :

// Inside message processing loop...

// Before leaseExpiration, for example, in a separate goroutine or timed check
// leaseToken is from the SELECT query
extendStmt := spanner.Statement{
    SQL: "SELECT * FROM RENEWLEASE_UserTasks(lease_tokens => [@token1])",
    Params: map[string]interface{}{
        "token1": leaseToken,
    },
}
_, err := client.Single().Query(ctx, extendStmt).Next() // Simplified call
if err != nil {
    log.Printf("Failed to extend lease for %s: %v", messageId, err)
    // Processing should probably stop as redelivery is likely
} else {
    // New lease expiration is typically approximately 10s from now
    log.Printf("Lease extended for %s", messageId)
    // Update local leaseExpiration time if needed
}

Accuser réception d'un message

Utilisez le LMD DELETE pour confirmer un message. Cette opération doit être transactionnelle avec toute autre écriture liée au traitement des messages.

GoogleSQL

-- DML for acknowledging
DELETE FROM UserTasks WHERE UserId = @userId AND MessageId = @messageId ASSERT_ROWS_MODIFIED 1;

PostgreSQL

-- DML for acknowledging
DELETE FROM usertasks WHERE userid = $1 AND messageid = $2 ASSERT_ROWS_MODIFIED 1;

Vous pouvez également utiliser la mutation Ack de la bibliothèque cliente pour confirmer un message :

Go

_, err := client.Apply(ctx, []*spanner.Mutation{
    spanner.Ack("UserTasks", spanner.Key{1}),
})

Java

dbClient.write(
    Collections.singletonList(
        Mutation.newAckBuilder("UserTasks")
            .setKey(Key.of(2L))
            .build()));

Bonnes pratiques

Voici quelques bonnes pratiques pour utiliser les files d'attente Spanner :

  • Charges utiles de petite taille : veillez à ce que les charges utiles des messages de file d'attente soient de petite taille (moins de 4 Ko). Utilisez le modèle de stockage hors bande pour les données plus volumineuses.
  • Gestion des baux : prolongez les baux pour les tâches qui peuvent dépasser le bail de message par défaut de 10 secondes. Si vous ne prolongez pas les baux des messages, cela peut entraîner des nouvelles distributions et un double traitement potentiel.
  • Gestion des erreurs : les files d'attente Spanner réessaient les messages dont le traitement a échoué avec un délai d'attente au cours de la première heure, et les messages plus anciens sont réessayés une fois par heure. Envisagez de déplacer les messages en échec permanent vers une file d'attente distincte.
  • Surveillance : surveillez la profondeur des files d'attente et l'âge des messages les plus anciens qui n'ont pas fait l'objet d'un accusé de réception pour vous assurer que vos récepteurs ne limitent pas l'ingestion de votre pipeline.
  • Idempotence : concevez vos processeurs de messages de manière à ce qu'ils tolèrent les opérations répétées, car la distribution de type "au moins une fois" signifie que des doublons occasionnels sont possibles. Pour en savoir plus, consultez la page Traitement de type "exactement une fois" et accusé de réception de type "au maximum une fois".
  • Durée de la fonction table : évitez les durées extrêmement courtes et excessivement longues pour les fonctions table. Nous vous recommandons une durée modérée, par exemple 20 minutes.
  • Taille du lot : ajustez max_batch_size à votre charge de travail. Utilisez des lots plus petits pour les événements à forte distribution ramifiée afin d'éviter les conflits de verrouillage sur les lignes partagées. Utilisez des lots plus importants pour les tâches indépendantes à nombre élevé de requêtes par seconde. Pour obtenir les meilleures performances, confirmez ou prolongez les baux de messages pour un lot en une seule transaction.

Surveiller

Vous pouvez surveiller les opérations de file d'attente à l'aide des tables d'introspection Spanner. Bien que ces tables n'incluent pas de colonnes spécifiques aux files d'attente, vous pouvez identifier l'activité des files d'attente en recherchant les noms de files d'attente définis par l'utilisateur dans les tables suivantes :

Les métriques de file d'attente Spanner suivantes se trouvent sous le préfixe spanner.googleapis.com/queue/* dans Cloud Monitoring :

  • buffered_ready_messages : (GAUGE, INT64, 1) Nombre de messages conservés en mémoire et prêts à être remis à un destinataire.
  • message_send_count : (DELTA, INT64, 1) nombre de messages envoyés dans Spanner au cours de l'intervalle pour une file d'attente.
  • message_ack_count : (DELTA, INT64, 1) nombre de messages reconnus dans Spanner au cours de l'intervalle pour une file d'attente.
  • oldest_unacked_message_age : (GAUGE, INT64, 1) Âge (en secondes) du message non confirmé le plus ancien d'une file d'attente.
  • lease_expiration_count : (DELTA, INT64, 1) Nombre d'expirations de bail dans Spanner au cours de l'intervalle pour une file d'attente.

Toutes les métriques précédentes sont échantillonnées environ toutes les 60 secondes. Après échantillonnage, les données ne sont pas visibles pendant un délai pouvant atteindre 120 secondes. La journalisation des audits Spanner couvre les opérations d'écriture, de lecture et de schéma sur les files d'attente.

Résoudre les problèmes

Les sections suivantes expliquent comment identifier et résoudre les problèmes courants lors de l'utilisation des files d'attente Spanner.

Le nombre de messages en attente augmente

Diagnostic

oldest_unacked_message_age et buffered_ready_messages sont tous deux élevés. Cela indique un déséquilibre entre votre taux d'envoi de messages et la capacité de traitement des messages de votre application.

Solution

Pour résoudre ce problème, procédez comme suit :

  • Vérifiez le taux d'accusé de réception : si la métrique message_ack_count a diminué, vérifiez que vos nœuds de calcul client s'exécutent correctement et qu'ils ne sont pas bloqués ni en panne.
  • Vérifiez le taux d'envoi : si message_send_count a augmenté, augmentez le nombre de nœuds de calcul de traitement des messages pour gérer la charge accrue.
  • Confirmez l'épuisement des ressources : vérifiez si les nombres buffered_ready_messages et lease_expiration_count sont élevés. Cette combinaison indique un manque de récepteurs de fonction de valeur de table (TVF) actifs ou un traitement client lent.

Des messages individuels sont bloqués

Diagnostic

La métrique oldest_unacked_message_age est élevée, mais buffered_ready_messages est faible ou stable. Cela indique que le traitement ou l'accusé de réception des messages individuels échouent, plutôt qu'un goulot d'étranglement de capacité global.

Solution

Pour résoudre ce problème, procédez comme suit :

  • Identifier les messages bloqués : interrogez la table de file d'attente et triez les résultats par heure de distribution pour trouver les messages non confirmés les plus anciens :

    GoogleSQL

    SELECT *
    FROM UserTasks
    ORDER BY DeliverTime ASC
    LIMIT 10;
    

    PostgreSQL

    SELECT *
    FROM usertasks
    ORDER BY deliver_time ASC
    LIMIT 10;
    
  • Enquêter sur les échecs de traitement : consultez les journaux de votre application pour déterminer pourquoi les nœuds de calcul ne reconnaissent pas les messages identifiés. Les messages non confirmés sont automatiquement renvoyés à l'expiration de leur bail.

Les baux de messages expirent avant la fin du traitement

Diagnostic

La métrique lease_expiration_count est élevée ou en augmentation. Cela indique que le temps de traitement des messages dépasse la durée du bail (10 secondes par défaut) avant que les nœuds de travail puissent accuser réception des messages.

Solution

Pour résoudre ce problème, procédez comme suit :

  • Renouvelez les baux de manière proactive : si le traitement des messages prend plus de 10 secondes, appelez la période RENEWLEASE_QUEUE_NAME() TVF de manière périodique. Renouvelez le bail environ sept à huit secondes après le début du traitement pour tenir compte de la latence du réseau.
  • Examinez le traitement lent : si votre application renouvelle déjà activement les baux, mais que lease_expiration_count reste élevé, vérifiez si votre code de backend présente des goulots d'étranglement de traitement, des appels RPC lents ou des blocages.

Impossible de faire évoluer le taux d'envoi des messages

Diagnostic

La métrique message_send_count atteint un plateau de débit ou les requêtes de publication rencontrent une latence d'écriture élevée lorsque vous tentez d'augmenter le taux d'envoi.

Solution

Pour résoudre ce problème, envisagez les modifications structurelles suivantes :

  • Augmentez les ressources de calcul : ajoutez des nœuds ou des unités de traitement à votre instance Spanner pour augmenter la capacité globale de la base de données.
  • Vérifiez les divisions : évaluez si vous pouvez ajouter d'autres divisions pour répartir la charge d'écriture sur plusieurs serveurs.
  • Optimisez le routage : assurez-vous que votre application de publication écrit directement dans la région principale de votre instance Spanner pour minimiser la latence d'écriture.

Les messages prêts ne sont pas traités

Diagnostic

La métrique buffered_ready_messages est élevée et augmente. Cela indique que les messages sont mis en mémoire tampon et prêts à être distribués en mémoire, mais que les nœuds de calcul du récepteur ne les extraient pas.

Solution

Pour résoudre ce problème, procédez comme suit :

  • Vérifiez les connexions TVF actives : vérifiez le nombre de TVF de récepteur actifs pour vous assurer que les workers du lecteur se connectent et extraient activement les messages. Si des nœuds de calcul sont déconnectés ou n'exécutent pas assez de requêtes RECEIVE_QUEUE_NAME() simultanées, les messages restent non interrogés dans le tampon.

Étapes suivantes