Spanner-Warteschlangen verwenden

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:

  1. Für nicht parsierbare Leasetokens wird keine Zeile zurückgegeben.
  2. Für bereits abgelaufene Leasetokens wird keine Zeile zurückgegeben.
  3. Für nicht erneuerbare Leasetokens wird eine Zeile mit einem SpannerNewLeaseToken von 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_size an 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:

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_count gesunken 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_count gestiegen 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_messages und lease_expiration_count erhö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_count aber 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