Utilizzare le code Spanner

Questo documento descrive come utilizzare le code Spanner. Spiega come creare una coda, inviare e ricevere messaggi, estendere i lease dei messaggi e confermare i messaggi. Include anche best practice, informazioni sulle code di monitoraggio e indicazioni per la risoluzione dei problemi.

Creare una coda

Per creare una coda, utilizza l'istruzione 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;

L'unica colonna che non deve essere creata esplicitamente nell'istruzione CREATE QUEUE si chiama DeliverTime in GoogleSQL e deliver_time in PostgreSQL. Vengono create automaticamente da Spanner.

Le code supportano le policy di durata (TTL), che possono aiutarti a gestire l'arretrato di messaggi vecchi e non riconosciuti.

Invia un messaggio

Per inviare un messaggio a una coda, utilizza l'istruzione INSERT DML:

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

In alternativa, utilizza le mutazioni Send della libreria client per inserire un messaggio:

Vai

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

Ricevere messaggi

Utilizza la funzione con valori di tabella (TVF) RECEIVE_QUEUE_NAME() con ExecuteStreamingSQL per ricevere messaggi. Questa è una chiamata di lunga durata. Devi eseguire una di queste chiamate per worker, per coda, in modo ciclico. Esegui la query utilizzando una lettura coerente perché Spanner rifiuta le letture obsolete.

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

Il codice client deve scorrere i risultati utilizzando una query di streaming. Ogni riga restituita è un messaggio.

Ricevere messaggi in batch

Per aumentare la velocità effettiva elaborando più messaggi insieme, puoi ricevere i messaggi in batch specificando l'argomento 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');

In alternativa, utilizza la libreria client. Questo esempio di Go mostra come trasmettere messaggi in streaming da una coda, verificare la scadenza del lease e riconoscere i messaggi in modo asincrono:

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

Estendere il lease del messaggio

Se l'elaborazione di un messaggio richiede più tempo del lease iniziale (10 secondi), utilizza la seguente sintassi:

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"

I token di lease vengono restituiti in base alla seguente logica:

  1. I token di lease non analizzabili non restituiscono una riga.
  2. I token di noleggio già scaduti non restituiscono una riga.
  3. I token di lease non rinnovabili restituiscono una riga con un SpannerNewLeaseToken di NULL. Ciò può verificarsi se il messaggio è già stato confermato, ma il token di lease non è scaduto.

In alternativa, utilizza la libreria client. Questo esempio di Go mostra come estendere la durata di un messaggio:

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

Confermare la ricezione di un messaggio

Utilizza l'istruzione DML DELETE per confermare un messaggio. Deve essere transazionale con qualsiasi altra scrittura correlata all'elaborazione dei messaggi.

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;

In alternativa, utilizza la mutazione Ack della libreria client per confermare la ricezione di un messaggio:

Vai

_, 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()));

Best practice

Di seguito sono riportate le best practice per l'utilizzo delle code Spanner:

  • Payload piccoli: mantieni i payload dei messaggi della coda di dimensioni inferiori a 4 kB. Utilizza il pattern di archiviazione out-of-band per i dati più grandi.
  • Gestione dei lease:estendi i lease per le attività che potrebbero superare il lease dei messaggi predefinito di 10 secondi. La mancata estensione dei lease dei messaggi può causare nuove consegne e una potenziale doppia elaborazione.
  • Gestione degli errori:le code Spanner riprovano a elaborare i messaggi non riusciti con backoff entro la prima ora, mentre i messaggi meno recenti verranno riprovati una volta all'ora. Valuta la possibilità di spostare i messaggi che non vengono recapitati in modo permanente in una coda separata.
  • Monitoraggio:monitora la profondità delle code e l'età del messaggio non riconosciuto più vecchio per assicurarti che i ricevitori non limitino l'acquisizione della pipeline.
  • Idempotenza:progetta i tuoi processori di messaggi in modo che tollerino le operazioni ripetute, perché la distribuzione "at-least-once" significa che è possibile che si verifichino duplicati occasionali. Per maggiori informazioni, consulta la pagina Elaborazione "exactly-once" e riconoscimento "at-most-once".
  • Durata della funzione con valori di tabella:evita durate estremamente brevi ed eccessivamente lunghe per le funzioni con valori di tabella. È consigliabile una durata moderata, ad esempio 20 minuti.
  • Dimensionamento batch:ottimizza max_batch_size in base al tuo carico di lavoro. Utilizza batch più piccoli per eventi con fan-out elevato per evitare conflitti di blocco sulle righe condivise. Utilizza batch più grandi per attività indipendenti con un numero elevato di query al secondo. Riconosci o estendi i lease dei messaggi per un batch in un'unica transazione per ottenere le migliori prestazioni.

Monitoraggio

Puoi monitorare le operazioni di coda utilizzando le tabelle di introspezione di Spanner. Anche se queste tabelle non includono colonne specifiche per le code, puoi identificare l'attività delle code cercando i nomi delle code definiti dall'utente nelle seguenti tabelle:

Le seguenti metriche della coda Spanner si trovano sotto il prefisso spanner.googleapis.com/queue/* in Cloud Monitoring:

  • buffered_ready_messages: (GAUGE, INT64, 1) Il conteggio dei messaggi memorizzati e pronti per la consegna a un destinatario.
  • message_send_count: (DELTA, INT64, 1) Il numero di messaggi inviati in Spanner durante l'intervallo per una coda.
  • message_ack_count: (DELTA, INT64, 1) Il numero di messaggi riconosciuti in Spanner durante l'intervallo per una coda.
  • oldest_unacked_message_age: (GAUGE, INT64, 1) L'età (in secondi) del messaggio più vecchio non riconosciuto in una coda.
  • lease_expiration_count: (DELTA, INT64, 1) Il numero di scadenze del lease in Spanner durante l'intervallo per una coda.

Tutte le metriche precedenti vengono campionate circa ogni 60 secondi. Dopo il campionamento, i dati potrebbero non essere visibili per un massimo di 120 secondi. L'audit logging di Spanner copre le operazioni di scrittura, lettura e schema sulle code.

Risoluzione dei problemi

Le sezioni seguenti descrivono come identificare e risolvere i problemi comuni quando utilizzi le code Spanner.

Il backlog dei messaggi è in aumento

Diagnosi

Sia oldest_unacked_message_age che buffered_ready_messages sono elevati. Ciò indica uno squilibrio tra la velocità di invio dei messaggi e la capacità di elaborazione dei messaggi della tua applicazione.

Risoluzione

Per risolvere il problema, segui questi passaggi:

  • Controlla il tasso di riconoscimento:se la metrica message_ack_count è diminuita, controlla i client worker per assicurarti che funzionino correttamente e non si siano bloccati o arrestati in modo anomalo.
  • Controlla la frequenza di invio:se message_send_count ha subito un picco, fai lo scale up dei worker di elaborazione dei messaggi per gestire il carico maggiore.
  • Conferma l'esaurimento delle risorse: controlla se i conteggi di buffered_ready_messages e lease_expiration_count sono elevati. Questa combinazione indica una carenza di ricevitori di funzioni con valori di tabella (TVF) attivi o un'elaborazione lenta del client.

I singoli messaggi sono bloccati

Diagnosi

La metrica oldest_unacked_message_age è elevata, ma buffered_ready_messages è bassa o stabile. Ciò indica che i singoli messaggi non vengono elaborati o riconosciuti, anziché un collo di bottiglia complessivo della capacità.

Risoluzione

Per risolvere il problema, segui questi passaggi:

  • Identifica i messaggi bloccati:esegui una query sulla tabella della coda e ordina i risultati in base ai tempi di consegna per trovare i messaggi non riconosciuti più vecchi:

    GoogleSQL

    SELECT *
    FROM UserTasks
    ORDER BY DeliverTime ASC
    LIMIT 10;
    

    PostgreSQL

    SELECT *
    FROM usertasks
    ORDER BY deliver_time ASC
    LIMIT 10;
    
  • Indaga sugli errori di elaborazione:controlla i log dell'applicazione per determinare perché i worker non riconoscono i messaggi identificati. I messaggi non riconosciuti vengono automaticamente riconsegnati dopo la scadenza del lease.

I contratti di locazione dei messaggi scadono prima del completamento dell'elaborazione

Diagnosi

La metrica lease_expiration_count è elevata o in aumento. Ciò indica che il tempo di elaborazione dei messaggi supera la durata del lease (10 secondi per impostazione predefinita) prima che i worker possano confermare i messaggi.

Risoluzione

Per risolvere il problema, segui questi passaggi:

  • Rinnova le concessioni in modo proattivo: se l'elaborazione del messaggio richiede più di 10 secondi, chiama periodicamente il TVF RENEWLEASE_QUEUE_NAME(). Rinnova il lease circa 7-8 secondi dopo l'inizio dell'elaborazione per tenere conto in modo sicuro della latenza di rete.
  • Analizza l'elaborazione lenta:se la tua applicazione rinnova già attivamente i lease, ma lease_expiration_count rimane elevato, controlla il codice di backend per individuare colli di bottiglia di elaborazione, chiamate RPC lente o deadlock.

Impossibile scalare la frequenza di invio dei messaggi

Diagnosi

La metrica message_send_count raggiunge un plateau di throughput o le richieste di pubblicazione riscontrano una latenza di scrittura elevata quando tenti di aumentare la velocità di invio.

Risoluzione

Per risolvere il problema, valuta le seguenti modifiche strutturali:

  • Fare lo scale up delle risorse di calcolo:aggiungi nodi o unità di elaborazione all'istanza Spanner per aumentare la capacità complessiva del database.
  • Controlla le suddivisioni:valuta se puoi aggiungere altre suddivisioni per distribuire il carico di scrittura su più server.
  • Ottimizza il routing:assicurati che l'applicazione di pubblicazione scriva direttamente nella regione leader dell'istanza Spanner per ridurre al minimo la latenza di scrittura.

I messaggi pronti non vengono elaborati

Diagnosi

La metrica buffered_ready_messages è elevata e in aumento. Ciò indica che i messaggi sono memorizzati nel buffer e pronti per la consegna in memoria, ma i worker riceventi non li recuperano.

Risoluzione

Per risolvere il problema, segui questi passaggi:

  • Controlla le connessioni TVF attive:controlla il numero di TVF ricevitore attivi per assicurarti che i worker di lettura si connettano e recuperino attivamente i messaggi. Se i worker si sono disconnessi o non eseguono un numero sufficiente di query RECEIVE_QUEUE_NAME() simultanee, i messaggi rimangono non sottoposti a polling nel buffer.

Passaggi successivi