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:
- I token di lease non analizzabili non restituiscono una riga.
- I token di noleggio già scaduti non restituiscono una riga.
- I token di lease non rinnovabili restituiscono una riga con un
SpannerNewLeaseTokendi 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_sizein 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:
- Statistiche sulle letture
- Statistiche sulle transazioni
- Statistiche sui blocchi
- Statistiche sulle dimensioni delle tabelle
- Statistiche sulle operazioni per tabella
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_countha 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_messageselease_expiration_countsono 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_countrimane 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
- Esplora altri scenari ed esempi di code Spanner.
- Scopri di più sull'elaborazione "exactly-once" e sull'acknowledgement "at-most-once".
- Configura il controllo dell'accesso con il controllo dell'controllo dell'accesso granulare per le code.