Szenarien und Beispiele für Spanner-Warteschlangen

In diesem Dokument finden Sie Architekturmuster und Codebeispiele für gängige Messaging-Szenarien mit Spanner-Warteschlangen. Mit diesen Mustern können Sie asynchrone Aufgaben 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 wurde, 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(NULL, NULL, '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 Standardleihfrist (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. Starten Sie die Berichterstellung.
  3. Rufen Sie SELECT * FROM RENEWLEASE_ReportQueue([leaseToken]) alle 5 Sekunden in einem separaten Thread oder einer separaten Routine auf.
  4. Bestätigen Sie nach Abschluss die Nachricht 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 (DELETE oder ACK) die aktuelle Warteschlangennachricht bei der Ankunft. Stellen Sie in derselben Transaktion eine neue Warteschlangennachricht mit einem Zustellungszeitstempel in der Zukunft in die Warteschlange, der über die für die Verarbeitung benötigte Zeit hinausgeht.
  2. Fahren Sie mit der Verarbeitung fort und bestätigen Sie die neu in die Warteschlange eingereihte 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 als Prüfpunkt speichern

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 Checkpointing-Muster verwenden:Die beste Methode für das Checkpointing 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: Für die Ausführung einer Bestellung sind Inventar und Zahlung erforderlich:

  1. Erstelle 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 Prüfung in die Warteschlange gestellt oder es wird eine andere Logik ausgeführt.

Regelmäßig eine Aktion ausführen

Wenn Sie eine Aktion regelmäßig ausführen möchten, verwenden Sie das Muster für die regelmäßige Planung. 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 kapselt:

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

Nächste Schritte