Szenarien und Beispiele für Spanner-Warteschlangen

Dieses Dokument enthält Architekturmuster und Codebeispiele für gängige Messaging-Szenarien mit Spanner-Warteschlangen. Mit diesen Mustern können Sie asynchrone Vorgänge nach dem Commit von Transaktionen auslösen, verzögerte oder wiederkehrende Aufgaben planen, große Nachrichtennutzlasten mit Out-of-Band-Speicher verwalten, Workflows mit mehreren Ereignissen koordinieren und Prüfpunkte oder Leases für lang andauernde Hintergrundjobs erstellen oder verlängern.

Genau einmalige Verarbeitung und höchstens einmalige Bestätigung

Die verschiedenen Überlegungen und Lösungen für die genau einmalige Verarbeitung und die höchstens einmalige Bestätigung werden auf der Seite Genau einmalige Verarbeitung und höchstens einmalige Bestätigung ausführlicher beschrieben.

Arbeit nach dem Commit einer Transaktion ausführen

Wenn Sie nach dem Commit einer Transaktion Arbeit ausführen möchten, senden Sie eine Nachricht an die Warteschlange innerhalb derselben Transaktion.

Wenn sich beispielsweise ein neuer Nutzer registriert, wird eine Willkommens-E-Mail gesendet:

GoogleSQL

-- Inside your application transaction:
-- 1. Insert into Users table
INSERT INTO Users (UserId, UserName) VALUES (124, 'New User');

-- 2. Send message to queue to trigger email
INSERT INTO UserTasks (UserId, MessageId, Payload)
VALUES (
  124,
  'welcome-email-id',
  b'{"type": "welcome", "email": "user@example.com"}'
);

PostgreSQL

-- Inside your application transaction:
-- 1. Insert into users table
INSERT INTO users (userid, username) VALUES (124, 'New User');

-- 2. Send message to queue to trigger email
INSERT INTO usertasks (userid, messageid, payload)
VALUES (
  124,
  'welcome-email-id',
  CAST('{"type": "welcome", "email": "user@example.com"}' AS bytea)
);

Nachdem die Transaktion abgeschlossen ist, streamt der Empfänger für UserTasks die Nachricht, sendet die E‑Mail und bestätigt die Nachricht:

GoogleSQL

-- 1. In the receiver process, stream messages from the queue
SELECT
  UserId,
  MessageId,
  Payload,
  DeliverTime,
  SpannerLeaseExpirationTimestamp,
  SpannerLeaseToken
FROM RECEIVE_UserTasks(max_duration=>'20m');

-- 2. After sending the welcome email, acknowledge the message
DELETE FROM UserTasks
WHERE UserId = 124 AND MessageId = 'welcome-email-id';

PostgreSQL

-- 1. In the receiver process, stream messages from the queue
SELECT
  userid,
  messageid,
  payload,
  deliver_time,
  spanner_lease_expiration_timestamp,
  spanner_lease_token
FROM spanner.receive_usertasks(
    max_batch_size=>NULL, priority=>NULL, max_duration=>'20m');

-- 2. After sending the welcome email, acknowledge the message
DELETE FROM usertasks
WHERE userid = 124 AND messageid = 'welcome-email-id';

Lang andauernde Aufgaben verarbeiten

Wenn Sie Aufgaben haben, die länger als die Standard-Lease-Zeit (mehr als 10 Sekunden) dauern, rufen Sie SELECT * FROM RENEWLEASE_QUEUE_NAME() regelmäßig auf.

Beispiel: einen Bericht erstellen:

  1. Der Empfänger erhält eine Nachricht von RECEIVE_ReportQueue().
  2. Berichterstellung starten
  3. Rufen Sie SELECT * FROM RENEWLEASE_ReportQueue([leaseToken]) alle 5 Sekunden in einem separaten Thread oder einer separaten Routine auf.
  4. Bestätigen Sie die Nachricht nach Abschluss und speichern Sie den Bericht.

Wenn Sie alternativ zeitaufwendige Aufgaben haben, die eine Verarbeitung vom Typ „At-most-once“ oder eine lange Lease-Zeit erfordern, gehen Sie so vor:

  1. Bestätigen Sie die aktuelle Warteschlangennachricht bei der Ankunft (DELETE oder ACK). Stellen Sie in derselben Transaktion eine neue Warteschlangennachricht mit einem Zustellungszeitstempel in der Zukunft in die Warteschlange, der über die Zeit hinausgeht, die für die Verarbeitung benötigt wird.
  2. Fahren Sie mit der Verarbeitung fort und bestätigen Sie die neu in die Warteschlange eingestellte Nachricht, wenn Sie fertig sind.

Die Vorteile dieses Ansatzes sind, dass die Lease nicht kontinuierlich verlängert werden muss und die Nachricht erst dann noch einmal gesendet wird, wenn der zukünftige Zeitpunkt erreicht ist (was Abstürze abdeckt). Wenn die erste Bestätigung erfolgreich ist, wird die höchstens einmalige Verarbeitung erreicht.

Lang andauernde Aufgaben mit Prüfpunkten versehen

Mit Spanner-Warteschlangen können Aufgaben verwaltet werden, die Minuten bis Stunden dauern, nicht nur schnelle Jobs. Verwenden Sie für diese zeitaufwendigen Aufgaben den folgenden Ansatz:

  1. Metadaten extern speichern:Verwenden Sie die Out-of-Band-Speicherung, um die Details und den Status der Aufgabe zu speichern.
  2. Regelmäßig Checkpoints erstellen:Um nach Abstürzen nicht viel Fortschritt zu verlieren, sollte der Status der Aufgabe regelmäßig gespeichert werden.
  3. Empfohlenes Prüfpunktmuster verwenden:Die beste Methode zum Erstellen von Prüfpunkten besteht darin, die aktuelle Warteschlangennachricht atomar zu bestätigen (ACK) und eine neue Nachricht zu senden, die für die zukünftige Zustellung geplant ist. Diese neue Nachricht enthält oder verweist auf den aktualisierten Status, wodurch eine sofortige erneute Zustellung an einen anderen Worker verhindert wird.

Dieses Muster reduziert die doppelte Arbeit, auch wenn kein vollständiges Checkpointing möglich ist. In diesem Fall wird die Aufgabe nach einem Absturz jedoch von Anfang an neu gestartet.

Aufgaben für einen bestimmten Zeitpunkt in der Zukunft planen

Wenn Sie Arbeit für einen bestimmten Zeitpunkt in der Zukunft planen möchten, legen Sie beim Einfügen der Nachricht die Spalte DeliverTime fest.

Beispiel für eine Erinnerung zum Ablauf des Testzeitraums:

GoogleSQL

-- 1. Insert into Users table
INSERT INTO Users (UserId, UserName) VALUES (125, 'Trial User');

-- 2. Send message to queue with a future delivery time
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (125, 'trial-expire-reminder', b'{"type": "reminder"}', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 29 DAY));

PostgreSQL

-- 1. Insert into users table
INSERT INTO users (userid, username) VALUES (125, 'Trial User');

-- 2. Send message to queue with a future delivery time
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (125, 'trial-expire-reminder', CAST('{"type": "reminder"}' AS bytea), CURRENT_TIMESTAMP + INTERVAL '29 DAY');

Große Nachrichtennutzlasten verarbeiten

Wenn die Nutzlast Ihrer Nachricht groß ist, verwenden Sie das Out-of-Band-Speichermuster. Speichern Sie die große Nutzlast in einer separaten Tabelle und fügen Sie einen Verweis darauf in die Warteschlangennachricht ein.

Beispiel: Bildverarbeitung

GoogleSQL

-- Schema
CREATE TABLE ImageUploads (
  UserId    INT64 NOT NULL,
  ImageId   STRING(36) NOT NULL,
  ImageData BYTES(MAX),
  Status    STRING(MAX) -- PENDING, PROCESSING, DONE
) PRIMARY KEY (UserId, ImageId),
  INTERLEAVE IN PARENT Users;

CREATE QUEUE ImageProcessingQueue (
  UserId    INT64 NOT NULL,
  ImageId   STRING(36) NOT NULL,
  Payload   BYTES(1) NOT NULL -- Payload can be minimal
) PRIMARY KEY (UserId, ImageId),
  INTERLEAVE IN PARENT ImageUploads ON DELETE CASCADE;

-- Application Logic
-- 1. Upload image, insert into ImageUploads with Status 'PENDING'
-- 2. Send message to ImageProcessingQueue
INSERT INTO ImageProcessingQueue (UserId, ImageId, Payload) VALUES (123, 'image-uuid-1', b'');

-- Receiver for ImageProcessingQueue:
-- 1. Receives message (UserId, ImageId).
-- 2. Reads ImageData from ImageUploads.
-- 3. Processes image.
-- 4. Updates ImageUploads Status to 'DONE'.
-- 5. ACKs the queue message.

PostgreSQL

-- Schema
CREATE TABLE imageuploads (
  userid    bigint NOT NULL,
  imageid   varchar(36) NOT NULL,
  imagedata bytea,
  status    varchar, -- PENDING, PROCESSING, DONE
  PRIMARY KEY (userid, imageid)
) INTERLEAVE IN PARENT users;

CREATE QUEUE imageprocessingqueue (
  userid    bigint NOT NULL,
  imageid   varchar(36) NOT NULL,
  payload   bytea NOT NULL, -- Payload can be minimal
  PRIMARY KEY (userid, imageid)
) INTERLEAVE IN PARENT imageuploads ON DELETE CASCADE;

-- Application Logic
-- 1. Upload image, insert into imageuploads with status 'PENDING'
-- 2. Send message to imageprocessingqueue
INSERT INTO imageprocessingqueue (userid, imageid, payload) VALUES (123, 'image-uuid-1', CAST('' AS bytea));

-- Receiver for imageprocessingqueue:
-- 1. Receives message (userid, imageid).
-- 2. Reads imagedata from imageuploads.
-- 3. Processes image.
-- 4. Updates imageuploads status to 'DONE'.
-- 5. ACKs the queue message.

Warten, bis mehrere Ereignisse eingetreten sind, bevor Sie fortfahren

Wenn Sie vor dem Fortfahren auf mehrere Ereignisse warten möchten (z. B. einen Join-Vorgang), verwenden Sie eine Tabelle zum Erfassen des Status und eine Warteschlange zum Auslösen von Prüfungen.

Beispiel: Auftragsabwicklung, für die Inventar und Zahlung erforderlich sind:

  1. Erstellen Sie eine Orders-Tabelle mit InventoryStatus und PaymentStatus.
  2. Wenn der Bestand bestätigt wurde, aktualisieren Sie Orders und senden Sie eine Nachricht an OrderCheckQueue.
  3. Wenn die Zahlung bestätigt wurde, aktualisiere Orders und sende eine Nachricht an OrderCheckQueue.
  4. Der Empfänger für OrderCheckQueue prüft die Tabelle Orders. Wenn beide Status bestätigt werden, wird der Versand fortgesetzt und die Nachricht wird bestätigt. Andernfalls wird sie möglicherweise für eine spätere Überprüfung in die Warteschlange eingereiht oder es wird eine andere Logik ausgeführt.

Aktion regelmäßig ausführen

Wenn Sie eine Aktion regelmäßig ausführen möchten, verwenden Sie das periodische Planungsmuster. Der Empfänger bestätigt die Nachricht und sendet eine neue Nachricht, die für das nächste Intervall geplant ist.

Beispiel für die stündliche Datenaggregation:

GoogleSQL

-- Inside your application transaction:
-- 1. Acknowledge current message
DELETE FROM AggregationQueue
WHERE TaskType = 'hourly-aggregator' AND MessageId = 'current-uuid'
ASSERT_ROWS_MODIFIED 1;

-- 2. Schedule next run 1 hour in the future
INSERT INTO AggregationQueue (TaskType, MessageId, Payload, DeliverTime)
VALUES ('hourly-aggregator', 'next-uuid', b'', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR));

PostgreSQL

-- Inside your application transaction:
-- 1. Acknowledge current message
DELETE FROM aggregationqueue
WHERE tasktype = 'hourly-aggregator' AND messageid = 'current-uuid'
ASSERT_ROWS_MODIFIED 1;

-- 2. Schedule next run 1 hour in the future
INSERT INTO aggregationqueue (tasktype, messageid, payload, deliver_time)
VALUES ('hourly-aggregator', 'next-uuid', CAST('' AS bytea), CURRENT_TIMESTAMP + INTERVAL '1 HOUR');

Alternativ können Sie die Clientbibliotheksmutationen Ack und Send verwenden. In diesen Beispielen wird davon ausgegangen, dass Sie ein Message-Objekt haben, das den Schlüssel und die Nutzlast enthält:

Java

// Receiver logic for AggregationQueue
public void process(DatabaseClient dbClient, Message msg) {
  // ... do aggregation ...

  // ACK current message and schedule next run (1 hour from now)
  Instant nextRun = Instant.now().plus(Duration.ofHours(1));
  Mutation ackMutation =
      Mutation.newAckBuilder("AggregationQueue")
          .setKey(msg.getKey()) // Ack
          .build();
  Mutation sendMutation =
      Mutation.newSendBuilder("AggregationQueue")
          .setKey(Key.of("hourly-aggregator", "next-uuid"))
          .setPayload(Value.bytes(ByteArray.copyFrom("")))
          .setDeliveryTime(nextRun) // Schedule next
          .build();
  dbClient.write(Arrays.asList(ackMutation, sendMutation));
}

Go

// Receiver logic for AggregationQueue
func process(msg) {
    // ... do aggregation ...

    // ACK current message and schedule next run
    nextRun := time.Now().Add(1 * time.Hour)
    _, err := client.Apply(ctx, []*spanner.Mutation{
        spanner.Ack("AggregationQueue", msg.Key), // Ack
        spanner.Send("AggregationQueue",
            spanner.Key{"hourly-aggregator", "next-uuid"},
            []byte(""),
            spanner.WithDeliveryTime(nextRun), // Schedule next
        ),
    })
    // ... handle err ...
}

Python

# Receiver logic for AggregationQueue
def process(database: spanner.Database, msg: Message):
  # ... do aggregation ...
  # ACK current message and schedule next run (1 hour from now)
  next_run = datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(
      hours=1
  )
  with database.batch() as batch:
    batch.ack(
        queue="AggregationQueue",
        key=msg.key,  # Ack
    )
    batch.send(
        queue="AggregationQueue",
        key=("hourly-aggregator", "next-uuid"),
        payload=b"",
        deliver_time=next_run,  # Schedule next
    )

Node.js

/**
 * Receiver logic for AggregationQueue
 * @param {import('@google-cloud/spanner').Database} database
 * @param { { key: Array<string|number>, payload: Buffer } } msg
 */
async function process(database, msg) {
  // ... do aggregation ...
  // ACK current message and schedule next run (1 hour from now)
  const nextRun = new Date(Date.now() + 60 * 60 * 1000);
  await database.runTransactionAsync(async (transaction) => {
    // Ack current message
    transaction.queueAck('AggregationQueue', msg.key);
    // Schedule next run
    transaction.queueSend(
      'AggregationQueue',
      ['hourly-aggregator', 'next-uuid'],
      {
        payload: Buffer.from(''),
        deliverTime: nextRun,
      }
    );
    await transaction.commit();
  });
}

Batch-Nachrichten mit temporärem Batching

Spanner-Warteschlangen stellen Nachrichten mit minimaler Latenz bereit. Wenn jedoch kontinuierlich eine große Anzahl von Nachrichten von unabhängigen Clients eingeht, kann die individuelle Verarbeitung jeder Nachricht zu einem hohen Transaktionsaufwand führen. Wenn Sie versuchen, die Warteschlangentabelle manuell abzufragen oder zu scannen, um Nachrichten in Batches zu verarbeiten, kann es zu Konflikten durch Bereichssperren, erhöhten Abbruchraten und zusätzlichen Lesekosten kommen.

Um Batching mit hohem Durchsatz ohne Konflikte zu erreichen, wenden Sie das temporale Batching-Muster an. Absender richten den DeliverTime von Nachrichten an einem diskreten Zeitfenster in der Zukunft aus (z. B. durch Aufrunden auf die nächste 10-Sekunden-Grenze). Da unabhängige Absender einen identischen zukünftigen Zeitstempel berechnen, werden die Nachrichten aus demselben Split in Spanner in der Regel gruppiert und in einem einzelnen Batch zugestellt, wenn max_batch_size in der tabellarischen Funktion RECEIVE_QUEUE_NAME() dies zulässt, oder in mehreren Batches, wenn max_batch_size kleiner als die Anzahl der zuzustellenden Nachrichten ist.

Der Zeitstempel für die Zustellung kann mit dieser Formel berechnet werden:

$$ \text{DeliverTime} = \text{RoundDown}(\text{CurrentTime}, \text{FixedDelay}) + \text{FixedDelay} $$

Bei einem 10‑Sekunden-Fenster erhalten beispielsweise alle Nachrichten, die zwischen 09:05:00 und 09:05:09.999 in die Warteschlange gestellt werden, den DeliverTime-Wert 09:05:10:

GoogleSQL

-- Calculate delivery time rounded to the next 10-second interval.
-- Use DIV(..., 10) * 10 to perform the RoundDown in SQL:
INSERT INTO OrderProcessingQueue (OrderId, DeliverTime, Payload)
VALUES (
  'order-101',
  TIMESTAMP_SECONDS(DIV(UNIX_SECONDS(CURRENT_TIMESTAMP()), 10) * 10 + 10),
  b'{"item": "book", "qty": 1}'
);

PostgreSQL

-- Calculate delivery time rounded to the next 10-second interval:
INSERT INTO orderprocessingqueue (orderid, deliver_time, payload)
VALUES (
  'order-101',
  to_timestamp((floor(extract(epoch from CURRENT_TIMESTAMP) / 10) * 10) + 10),
  CAST('{"item": "book", "qty": 1}' AS bytea)
);

Empfänger rufen diese zeitlich abgestimmten Nachrichten dann in Batches mit max_batch_size ab:

GoogleSQL

SELECT OrderId, Payload, DeliverTime, SpannerLeaseToken
FROM RECEIVE_OrderProcessingQueue(max_duration=>'20m', max_batch_size=>50);

PostgreSQL

SELECT orderid, payload, deliver_time, spanner_lease_token
FROM spanner.receive_orderprocessingqueue(
    max_batch_size=>50, priority=>NULL, max_duration=>'20m');

Wenn Sie Millionen von Nachrichten gleichzeitig senden, kann es zu plötzlichen Verarbeitungsspitzen kommen, wenn alle Nachrichten auf die exakt gleiche Sekunde ausgerichtet werden. Wenn Sie die Arbeit gleichmäßig verteilen und Nachrichten trotzdem nach Einheit gruppieren möchten, fügen Sie der Berechnung einen Offset basierend auf einer eindeutigen Kennung (z. B. TIMESTAMP_ADD(deliver_time, INTERVAL MOD(ABS(FARM_FINGERPRINT(CAST(UserId AS STRING))), 10) SECOND)) hinzu.

Datenänderung von der Verarbeitung entkoppeln (Dirty-Flag-Muster)

In transaktionalen Anwendungen mit hohem Durchsatz kann die Ausführung komplexer Neuberechnungen, Suchindexierungen oder Cache-Invalidierungen direkt in nutzerorientierten Transaktionen die Nutzerfreundlichkeit beeinträchtigen, da die Latenz erhöht und die Sperrkonflikte in gemeinsam genutzten Zeilen verursacht werden.

Das Dirty-Flag-Muster entkoppelt Datenänderungen von der asynchronen Verarbeitung. Wenn eine Transaktion eine Tabelle ändert, wird innerhalb derselben Transaktion eine einfache „Dirty Bit“-Nachricht in eine Warteschlange geschrieben. Ein Hintergrundworker verarbeitet die Nachricht dann und führt die aufwendige Verarbeitung asynchron aus.

Wenn ein Kunde beispielsweise sein Profil oder seine Einstellungen ändert:

GoogleSQL

-- Inside user profile update transaction:
-- 1. Update the primary entity table
UPDATE UserProfiles
SET FullName = 'Jane Doe', UpdatedAt = CURRENT_TIMESTAMP()
WHERE UserId = 456;

-- 2. Send lightweight dirty flag message to the queue.
-- It is recommended to interleave the queue in the primary UserProfiles
-- table for better transaction performance.
INSERT INTO UserDirtyQueue (UserId, TaskType, CommitTimestamp, Payload)
VALUES (456, 'reindex-user-profile', CURRENT_TIMESTAMP(), b'');

PostgreSQL

-- Inside user profile update transaction:
-- 1. Update the primary entity table
UPDATE userprofiles
SET fullname = 'Jane Doe', updatedat = CURRENT_TIMESTAMP
WHERE userid = 456;

-- 2. Send lightweight dirty flag message to the queue
INSERT INTO userdirtyqueue (userid, tasktype, committimestamp, payload)
VALUES (456, 'reindex-user-profile', CURRENT_TIMESTAMP, CAST('' AS bytea));

Der Hintergrund-Receiver für UserDirtyQueue empfängt UserId, liest die neue Profilzeile außerhalb des kritischen Pfads des Nutzers und berechnet den Suchindex neu oder aktualisiert externe Caches. Batching kann auch auf dieses Muster angewendet werden, wenn mehrere Aktualisierungen für denselben Nutzer innerhalb kurzer Zeit an die Warteschlange gesendet werden. In diesem Fall kann der TVF-Empfänger einen max_batch_size größer als 1 angeben, um mehrere Nachrichten aus demselben Batch zu empfangen.

Worker-Zustand überwachen und Zeitüberschreitungen erkennen

Sie können geplante Warteschlangennachrichten verwenden, um ein fehlertolerantes Heartbeat- und Integritätsüberwachungssystem für Flotten von Worker-Knoten oder Mikrodienstinstanzen zu erstellen.

So implementieren Sie Systemdiagnosen:

  1. Worker bei Start registrieren:Wenn ein Worker initialisiert wird, fügt er eine Heartbeat-Nachricht in eine Systemdiagnosewarteschlange ein, wobei DeliverTime auf die Frist für Fehler (z. B. 60 Sekunden) festgelegt wird.
  2. Regelmäßige Heartbeats senden:Solange der Worker fehlerfrei funktioniert, aktualisiert er seine Heartbeat-Nachricht regelmäßig (z. B. alle 10 Sekunden), indem er DeliverTime um weitere 60 Sekunden in die Zukunft verschiebt.
  3. Fehler erkennen:Wenn der Worker abstürzt oder die Netzwerkverbindung verliert, werden keine Heartbeats mehr gesendet. Nach 60 Sekunden wird der Zustellungszeitstempel aktualisiert (DeliverTime <= CURRENT_TIMESTAMP()) und Spanner sendet die Nachricht an einen Empfänger für Benachrichtigungen, der ein Failover oder eine Neuzuweisung der Aufgabe initiiert.

Wichtig:Spanner-Warteschlangen unterstützen keine UPDATE-DML-Anweisungen. Wenn Sie den Zeitstempel für den Heartbeat aktualisieren möchten, müssen Sie die vorhandene Nachricht löschen und in einer einzelnen Transaktion eine Ersatznachricht mit dem neuen DeliverTime einfügen oder die Mutationen Ack und Send der Clientbibliothek anwenden. Achten Sie darauf, dass der Primärschlüssel der Warteschlange nur WorkerId (und nicht (WorkerId, MessageId)) ist, damit zu jedem Zeitpunkt nur eine Heartbeat-Nachricht pro Worker vorhanden ist.

GoogleSQL

-- Inside the worker heartbeat transaction (executed every 10 seconds):
-- 1. Acknowledge the existing heartbeat message
DELETE FROM WorkerHealthQueue
WHERE WorkerId = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

-- 2. Send replacement heartbeat with refreshed 60-second deadline
INSERT INTO WorkerHealthQueue (WorkerId, DeliverTime, Payload)
VALUES (
  'worker-node-42',
  TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 60 SECOND),
  b'{"status": "healthy", "active_jobs": 3}'
);

PostgreSQL

-- Inside the worker heartbeat transaction (executed every 10 seconds):
-- 1. Acknowledge the existing heartbeat message
DELETE FROM workerhealthqueue
WHERE workerid = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

-- 2. Send replacement heartbeat with refreshed 60-second deadline
INSERT INTO workerhealthqueue (workerid, deliver_time, payload)
VALUES (
  'worker-node-42',
  CURRENT_TIMESTAMP + INTERVAL '60 SECOND',
  CAST('{"status": "healthy", "active_jobs": 3}' AS bytea)
);

Wenn Sie Clientbibliotheksmutationen verwenden, muss Ack vor Send im Mutations-Slice stehen, wie in diesem Go-Beispiel:

// Worker heartbeat loop
func sendHeartbeat(ctx context.Context, client *spanner.Client, workerID string) error {
    newDeadline := time.Now().Add(60 * time.Second)
    _, err := client.Apply(ctx, []*spanner.Mutation{
        spanner.Ack("WorkerHealthQueue", spanner.Key{workerID}),
        spanner.Send(
            "WorkerHealthQueue",
            spanner.Key{workerID},
            []byte(`{"status":"healthy"}`),
            spanner.WithDeliveryTime(newDeadline),
        ),
    })
    return err
}

Wenn ein Worker ordnungsgemäß heruntergefahren wird, wird seine Heartbeat-Nachricht explizit gelöscht, damit kein falscher Alarm ausgelöst wird:

DELETE FROM WorkerHealthQueue WHERE WorkerId = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

Mehrere Aufgabentypen in einer einzelnen Warteschlange verarbeiten

Für Spanner-Instanzen gelten Limits für die Gesamtzahl der Warteschlangen. Wenn Sie für jeden kleinen asynchronen Vorgang eine separate Warteschlange erstellen, kann dieses Limit schnell erreicht werden. Außerdem müssen Sie dann viele gleichzeitige Empfängeranfragen verwalten.

Um Vorgänge zu konsolidieren, können Sie verschiedene Aufgabentypen in einer einzigen Warteschlange kombinieren, die als polymorphe Warteschlange bezeichnet wird. Es gibt zwei Strategien zum Erstellen einer polymorphen Warteschlange.

Strategie 1: Typspalte in den Primärschlüssel aufnehmen

GoogleSQL

CREATE QUEUE ApplicationTasks (
  TaskType   STRING(50) NOT NULL,
  TaskId     STRING(36) NOT NULL,
  Payload    BYTES(MAX) NOT NULL,
) PRIMARY KEY (TaskType, TaskId);

-- Enqueue an email task
INSERT INTO ApplicationTasks (TaskType, TaskId, Payload)
VALUES ('SEND_EMAIL', 'task-uuid-1', b'{"to": "user@example.com", "template": "welcome"}');

-- Enqueue an image thumbnail task
INSERT INTO ApplicationTasks (TaskType, TaskId, Payload)
VALUES ('GENERATE_THUMBNAIL', 'task-uuid-2', b'{"image_id": "img-789", "size": "small"}');

PostgreSQL

CREATE QUEUE applicationtasks (
  tasktype   varchar(50) NOT NULL,
  taskid     varchar(36) NOT NULL,
  payload    bytea NOT NULL,
  PRIMARY KEY (tasktype, taskid)
);

-- Enqueue an email task
INSERT INTO applicationtasks (tasktype, taskid, payload)
VALUES ('SEND_EMAIL', 'task-uuid-1', CAST('{"to": "user@example.com", "template": "welcome"}' AS bytea));

-- Enqueue an image thumbnail task
INSERT INTO applicationtasks (tasktype, taskid, payload)
VALUES ('GENERATE_THUMBNAIL', 'task-uuid-2', CAST('{"image_id": "img-789", "size": "small"}' AS bytea));

Der Empfänger prüft TaskType und leitet die Nutzlast an den entsprechenden Handler weiter.

Strategie 2: Polymorphe Nutzlaststruktur

Alternativ können Sie eine JSON-Nutzlast mit einem Feld für die Aktions- oder Typunterscheidung verwenden:

{
  "action": "SYNC_INVENTORY",
  "data": { "item_id": 987, "delta": -1 }
}

Benutzerdefinierte Verzögerungen für Wiederholungsversuche implementieren

In Spanner-Warteschlangen werden fehlgeschlagene oder nicht bestätigte Nachrichten automatisch mit integriertem exponentiellem Backoff noch einmal gesendet. In Szenarien, in denen eine Nachricht aus einem bekannten Grund mit einer bekannten Dauer fehlschlägt oder ein externes Ratenlimit (z. B. eine HTTP-429-Antwort mit einem Retry-After-Header) vorliegt, kann das automatische Backoff zu vorzeitigen Wiederholungsversuchen führen, die CPU-Ressourcen verschwenden.

So implementieren Sie eine benutzerdefinierte Verzögerung für Wiederholungsversuche:

  1. Fangen Sie den spezifischen vorübergehenden Fehler in Ihrem Message Processor ab.
  2. Bestätigen Sie die aktuelle Nachricht, um den aktuellen Zustellversuch abzuschließen.
  3. Senden Sie in derselben Transaktion eine Ersatznachricht mit einem expliziten DeliverTime, das auf die ausgewählte zukünftige Wiederholungszeit festgelegt ist (im folgenden Beispiel 5 Minuten später).

GoogleSQL

-- Inside failure-handling transaction:
-- 1. Acknowledge the failed message
DELETE FROM OutboundNotificationQueue
WHERE NotificationId = 'notif-555' ASSERT_ROWS_MODIFIED 1;

-- 2. Reschedule delivery 5 minutes in the future
INSERT INTO OutboundNotificationQueue (NotificationId, DeliverTime, Payload)
VALUES (
  'notif-555',
  TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 5 MINUTE),
  b'{"recipient": "user@example.com", "retry_count": 2}'
);

PostgreSQL

-- Inside failure-handling transaction:
-- 1. Acknowledge the failed message
DELETE FROM outboundnotificationqueue
WHERE notificationid = 'notif-555' ASSERT_ROWS_MODIFIED 1;

-- 2. Reschedule delivery 5 minutes in the future
INSERT INTO outboundnotificationqueue (notificationid, deliver_time, payload)
VALUES (
  'notif-555',
  CURRENT_TIMESTAMP + INTERVAL '5 MINUTE',
  CAST('{"recipient": "user@example.com", "retry_count": 2}' AS bytea)
);

Mit diesem Ansatz kann Ihre Anwendung Backoff-Zeitpläne genau verwalten und eine Überlastung externer APIs während der Downstream-Wiederherstellungsphasen vermeiden.

Nächste Schritte