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:
- Der Empfänger erhält eine Nachricht von
RECEIVE_ReportQueue(). - Starten Sie die Berichterstellung.
- Rufen Sie
SELECT * FROM RENEWLEASE_ReportQueue([leaseToken])alle 5 Sekunden in einem separaten Thread oder einer separaten Routine auf. - 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:
- Bestätigen Sie (
DELETEoderACK) 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. - 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:
- Metadaten extern speichern:Verwenden Sie die Out-of-Band-Speicherung, um die Details und den Status der Aufgabe zu speichern.
- Regelmäßig Checkpoints erstellen:Um nach Abstürzen nicht viel Fortschritt zu verlieren, sollte der Status der Aufgabe regelmäßig gespeichert werden.
- 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:
- Erstelle eine
Orders-Tabelle mitInventoryStatusundPaymentStatus. - Wenn der Bestand bestätigt wurde, aktualisieren Sie
Ordersund senden Sie eine Nachricht anOrderCheckQueue. - Wenn die Zahlung bestätigt wurde, aktualisiere
Ordersund sende eine Nachricht anOrderCheckQueue. - Der Empfänger für
OrderCheckQueueprüft die TabelleOrders. 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
- Informationen zur Verwendung von Spanner-Warteschlangen, einschließlich Best Practices und Monitoring
- Genau einmalige Verarbeitung und höchstens einmalige Bestätigung
- Konfigurieren Sie die Zugriffssteuerung mit der detaillierten Zugriffssteuerung für Warteschlangen.