In diesem Dokument wird beschrieben, wie Sie Spanner-Warteschlangen verwenden. Darin wird beschrieben, wie Sie eine Warteschlange erstellen, Nachrichten senden und empfangen, Nachrichtenleases verlängern und Nachrichten bestätigen. Außerdem enthält es Best Practices, Informationen zum Überwachen von Warteschlangen und Anleitungen zur Fehlerbehebung.
Warteschlange erstellen
Verwenden Sie die CREATE QUEUE-Anweisung, um eine Warteschlange zu erstellen.
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;
Die einzige Spalte, die in der CREATE
QUEUE-Anweisung nicht explizit erstellt werden muss, heißt in GoogleSQL DeliverTime und in PostgreSQL deliver_time. Sie werden automatisch von Spanner erstellt.
Warteschlangen unterstützen Richtlinien zur Gültigkeitsdauer (TTL), mit denen sich der Nachrichtenrückstand für alte, nicht bestätigte Nachrichten verwalten lässt.
Nachricht senden
Verwenden Sie zum Senden einer Nachricht an eine Warteschlange die DML-Anweisung 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');
Alternativ können Sie die Clientbibliothek Send-Mutationen zum Einfügen einer Nachricht verwenden:
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()));
Nachrichten empfangen
Verwenden Sie die Tabellenwertfunktion RECEIVE_QUEUE_NAME() (Table-valued function, TVF) mit ExecuteStreamingSQL, um Nachrichten zu empfangen. Dies ist ein Anruf mit langer Dauer. Sie müssen einen dieser Aufrufe pro Worker und pro Warteschlange in einer Schleife ausführen. Führen Sie die Abfrage mit einem starken Lesevorgang aus, da Spanner veraltete Lesevorgänge ablehnt.
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');
Der Clientcode sollte die Ergebnisse mit einer Streaming-Abfrage durchlaufen. Jede zurückgegebene Zeile ist eine Nachricht.
Nachrichten in einem Batch empfangen
Wenn Sie den Durchsatz erhöhen möchten, indem Sie mehrere Nachrichten zusammen verarbeiten, können Sie Nachrichten in Batches empfangen. Geben Sie dazu das Argument max_batch_size an:
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');
Alternativ können Sie die Clientbibliothek verwenden. In diesem Go-Beispiel wird gezeigt, wie Nachrichten aus einer Warteschlange gestreamt, der Ablauf von Leases überprüft und Nachrichten asynchron bestätigt werden:
// 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)
}
Nachricht-Lease verlängern
Wenn die Verarbeitung einer Nachricht länger als die ursprüngliche Leasedauer (10 Sekunden) dauert, verwenden Sie die folgende Syntax:
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"
Lease-Tokens werden gemäß der folgenden Logik zurückgegeben:
- Für nicht parsierbare Leasetokens wird keine Zeile zurückgegeben.
- Für bereits abgelaufene Leasetokens wird keine Zeile zurückgegeben.
- Für nicht erneuerbare Leasetokens wird eine Zeile mit einem
SpannerNewLeaseTokenvon NULL zurückgegeben. Das kann passieren, wenn die Nachricht bereits bestätigt wurde, das Lease-Token aber noch nicht abgelaufen ist.
Alternativ können Sie die Clientbibliothek verwenden. In diesem Go-Beispiel wird gezeigt, wie Sie eine Nachrichtenlease verlängern:
// 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
}
Nachricht bestätigen
Verwenden Sie die DELETE-DML, um eine Nachricht zu bestätigen. Dies muss mit allen anderen Schreibvorgängen im Zusammenhang mit der Nachrichtenverarbeitung transaktional sein.
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;
Alternativ können Sie die Clientbibliotheksmutation Ack verwenden, um eine Nachricht zu bestätigen:
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()));
Best Practices
Im Folgenden finden Sie Best Practices für die Verwendung von Spanner-Warteschlangen:
- Kleine Nutzlasten:Die Nutzlasten von Warteschlangennachrichten sollten unter 4 KB bleiben. Verwenden Sie das Out-of-Band-Speichermuster für größere Datenmengen.
- Freigabeverwaltung:Verlängern Sie die Freigaben für Aufgaben, die die standardmäßige Nachrichtenfreigabe von 10 Sekunden überschreiten könnten. Wenn Nachrichtenleihfristen nicht verlängert werden, kann dies zu erneuten Zustellungen und einer potenziellen doppelten Verarbeitung führen.
- Fehlerbehandlung:Spanner-Warteschlangen versuchen, Nachrichten, bei denen die Verarbeitung innerhalb der ersten Stunde fehlgeschlagen ist, mit Backoff-Verfahren noch einmal zu verarbeiten. Ältere Nachrichten werden einmal pro Stunde noch einmal verarbeitet. Verschieben Sie Nachrichten, die dauerhaft fehlschlagen, in eine separate Warteschlange.
- Monitoring:Überwachen Sie die Warteschlangentiefe und das Alter der ältesten nicht bestätigten Nachricht, um sicherzustellen, dass Ihre Empfänger die Aufnahme in die Pipeline nicht einschränken.
- Idempotenz:Entwerfen Sie Ihre Nachrichtenprozessoren so, dass sie wiederholte Vorgänge tolerieren, da bei der mindestens einmaligen Zustellung gelegentliche Duplikate möglich sind. Weitere Informationen finden Sie auf der Seite Exactly-once processing and at-most-once acknowledgment.
- Dauer der Tabellenwertfunktion:Vermeiden Sie sowohl extrem kurze als auch extrem lange Dauern für Tabellenwertfunktionen. Wir empfehlen eine moderate Dauer, z. B. 20 Minuten.
- Batchgröße:Passen Sie
max_batch_sizean Ihre Arbeitslast an. Verwenden Sie kleinere Batches für Ereignisse mit hohem Fan-out, um Sperrkonflikte bei gemeinsam genutzten Zeilen zu vermeiden. Verwenden Sie größere Batches für unabhängige Aufgaben mit vielen Abfragen pro Sekunde. Bestätigen oder verlängern Sie Nachrichtenleases für einen Batch in einer einzigen Transaktion, um die beste Leistung zu erzielen.
Überwachen
Sie können Warteschlangenoperationen mit Spanner-Introspection-Tabellen überwachen. Diese Tabellen enthalten zwar keine warteschlangenspezifischen Spalten, Sie können die Warteschlangenaktivität aber ermitteln, indem Sie in den folgenden Tabellen nach Ihren benutzerdefinierten Warteschlangennamen suchen:
- Statistiken lesen
- Transaktionsstatistiken
- Statistiken sperren
- Statistiken zur Tabellengröße
- Statistiken zu Tabellenvorgängen
Die folgenden Spanner-Warteschlangenmesswerte finden Sie in Cloud Monitoring unter dem Präfix spanner.googleapis.com/queue/*:
buffered_ready_messages: (GAUGE, INT64, 1) Die Anzahl der Nachrichten, die im Arbeitsspeicher gehalten werden und für die Zustellung an einen Empfänger bereit sind.message_send_count: (DELTA, INT64, 1) Die Anzahl der Nachrichten, die während des Intervalls für eine Warteschlange in Spanner gesendet wurden.message_ack_count: (DELTA, INT64, 1) Die Anzahl der Nachrichten, die in Spanner während des Intervalls für eine Warteschlange bestätigt wurden.oldest_unacked_message_age: (GAUGE, INT64, 1) Das Alter (in Sekunden) der ältesten nicht bestätigten Nachricht in einer Warteschlange.lease_expiration_count: (DELTA, INT64, 1) Die Anzahl der Lease-Abläufe in Spanner während des Intervalls für eine Warteschlange.
Alle oben genannten Messwerte werden etwa alle 60 Sekunden abgerufen. Nach dem Abruf werden bis zu 120 Sekunden lang möglicherweise keine Daten angezeigt. Spanner-Audit-Logging umfasst Schreib-, Lese- und Schemavorgänge für Warteschlangen.
Fehlerbehebung
In den folgenden Abschnitten wird beschrieben, wie Sie häufige Probleme bei der Verwendung von Spanner-Warteschlangen erkennen und beheben.
Anzahl der ausstehenden Nachrichten nimmt zu
Diagnose
Sowohl oldest_unacked_message_age als auch buffered_ready_messages sind erhöht. Dies weist auf ein Ungleichgewicht zwischen der Rate, mit der Sie Nachrichten senden, und der Kapazität Ihrer Anwendung zur Verarbeitung von Nachrichten hin.
Lösung
So beheben Sie das Problem:
- Bestätigungsrate prüfen:Wenn der Messwert
message_ack_countgesunken ist, prüfen Sie Ihre Client-Worker, um sicherzustellen, dass sie ordnungsgemäß ausgeführt werden und nicht hängen geblieben oder abgestürzt sind. - Sendegeschwindigkeit prüfen:Wenn
message_send_countgestiegen ist, skalieren Sie die Worker für die Nachrichtenverarbeitung hoch, um die erhöhte Last zu bewältigen. - Ressourcenerschöpfung bestätigen:Prüfen Sie, ob die Anzahl der
buffered_ready_messagesundlease_expiration_counterhöht ist. Diese Kombination deutet auf einen Mangel an aktiven Empfängern von Tabellenwertfunktionen (Table-valued Function, TVF) oder eine langsame Clientverarbeitung hin.
Einzelne Nachrichten bleiben hängen
Diagnose
Der Messwert oldest_unacked_message_age ist hoch, aber buffered_ready_messages ist niedrig oder stabil. Das bedeutet, dass einzelne Nachrichten nicht verarbeitet oder bestätigt werden können, und nicht, dass es ein allgemeines Kapazitätsengpass gibt.
Lösung
So beheben Sie das Problem:
Hängende Nachrichten identifizieren:Fragen Sie die Warteschlangentabelle ab und sortieren Sie die Ergebnisse nach der Zustellzeit, um die ältesten nicht bestätigten Nachrichten zu finden:
GoogleSQL
SELECT * FROM UserTasks ORDER BY DeliverTime ASC LIMIT 10;PostgreSQL
SELECT * FROM usertasks ORDER BY deliver_time ASC LIMIT 10;Verarbeitungsfehler untersuchen:Sehen Sie in den Anwendungslogs nach, warum die Worker die identifizierten Nachrichten nicht bestätigen. Unbestätigte Nachrichten werden nach Ablauf ihrer Lease automatisch noch einmal zugestellt.
Nachrichtensperren laufen ab, bevor die Verarbeitung abgeschlossen ist
Diagnose
Der Messwert lease_expiration_count ist erhöht oder nimmt zu. Dies weist darauf hin, dass die Nachrichtenverarbeitungszeit die Lease-Dauer (standardmäßig 10 Sekunden) überschreitet, bevor Worker die Nachrichten bestätigen können.
Lösung
So beheben Sie das Problem:
- Leases proaktiv verlängern:Wenn die Verarbeitung von Nachrichten länger als 10 Sekunden dauert, rufen Sie das
RENEWLEASE_QUEUE_NAME()TVF regelmäßig auf. Erneuern Sie die Zuweisung etwa 7–8 Sekunden nach Beginn der Verarbeitung, um Netzwerk-Latenzzeiten zu berücksichtigen. - Langsame Verarbeitung untersuchen:Wenn deine Anwendung bereits aktiv Leases erneuert,
lease_expiration_countaber weiterhin hoch ist, prüfe deinen Backend-Code auf Verarbeitungsengpässe, langsame RPC-Aufrufe oder Deadlocks.
Nachrichtensenderate kann nicht skaliert werden
Diagnose
Der Messwert message_send_count erreicht ein Durchsatzplateau oder bei Veröffentlichungsanfragen tritt eine erhöhte Schreiblatenz auf, wenn Sie versuchen, die Senderate zu erhöhen.
Lösung
Um dieses Problem zu beheben, sollten Sie die folgenden strukturellen Änderungen in Betracht ziehen:
- Rechenressourcen vertikal skalieren:Fügen Sie Ihrer Spanner-Instanz Knoten oder Verarbeitungseinheiten hinzu, um die Gesamtdatenbankkapazität zu erhöhen.
- Splits prüfen:Untersuchen Sie, ob Sie weitere Splits hinzufügen können, um die Schreiblast auf mehrere Server zu verteilen.
- Routing optimieren:Achten Sie darauf, dass Ihre Veröffentlichungsanwendung direkt in die führende Region Ihrer Spanner-Instanz schreibt, um die Schreiblatenz zu minimieren.
„Bereit“-Nachrichten werden nicht verarbeitet
Diagnose
Der Messwert buffered_ready_messages ist hoch und nimmt zu. Das bedeutet, dass Nachrichten im Arbeitsspeicher gepuffert und für die Zustellung bereit sind, aber nicht von den Empfänger-Workern abgerufen werden.
Lösung
So beheben Sie das Problem:
- Aktive TVF-Verbindungen prüfen:Prüfen Sie die Anzahl der aktiven TVF-Empfänger, um sicherzustellen, dass Reader-Worker aktiv eine Verbindung herstellen und Nachrichten abrufen. Wenn die Worker die Verbindung getrennt haben oder nicht genügend gleichzeitige
RECEIVE_QUEUE_NAME()-Anfragen ausführen, bleiben Nachrichten ungeprüft im Puffer.
Nächste Schritte
- Weitere Szenarien und Beispiele für Cloud Spanner-Warteschlangen
- Genau einmalige Verarbeitung und höchstens einmalige Bestätigung
- Konfigurieren Sie die Zugriffssteuerung mit der detaillierten Zugriffssteuerung für Warteschlangen.