Ce document fournit des modèles architecturaux et des exemples de code pour les scénarios de messagerie courants utilisant des files d'attente Spanner. Vous pouvez utiliser ces modèles pour déclencher des tâches asynchrones après la validation des transactions, planifier des tâches différées ou récurrentes, gérer de grandes charges utiles de messages avec le stockage hors bande, coordonner des workflows multi-événements, et créer des points de contrôle ou étendre des baux pour les tâches de longue durée en arrière-plan.
Traitement de type "exactement une fois" et accusé de réception "au maximum une fois"
Les différentes considérations et solutions pour le traitement de type "exactement une fois" et l'accusé de réception de type "au plus une fois" sont décrites plus en détail sur la page Traitement de type "exactement une fois" et accusé de réception de type "au plus une fois".
Effectuer des tâches après la validation d'une transaction
Pour effectuer des tâches après la validation d'une transaction, envoyez un message à la file d'attente dans la même transaction.
Par exemple, l'inscription d'un nouvel utilisateur déclenche l'envoi d'un e-mail de bienvenue :
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)
);
Une fois la transaction validée, le récepteur de UserTasks diffuse le message, envoie l'e-mail et accuse réception du message :
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';
Gérer les tâches de longue durée
Si vous avez une tâche qui peut prendre plus de temps que la durée de bail par défaut (plus de 10 secondes), appelez SELECT * FROM RENEWLEASE_QUEUE_NAME() de manière périodique.
Par exemple, pour générer un rapport :
- Le destinataire reçoit un message de
RECEIVE_ReportQueue(). - Lancez la génération du rapport.
- Appelez
SELECT * FROM RENEWLEASE_ReportQueue([leaseToken])toutes les cinq secondes dans un thread ou une routine distincts. - Une fois l'opération terminée, accusez réception du message et stockez le rapport.
Si vous avez un travail de longue durée qui nécessite un traitement "au plus une fois" ou une longue durée de bail, procédez comme suit :
- Confirmez (
DELETEouACK) le message actuel de la file d'attente à l'arrivée. Dans la même transaction, remettez en file d'attente un nouveau message avec un code temporel de remise futur, au-delà du temps nécessaire au traitement. - Poursuivez le traitement et confirmez le nouveau message mis en file d'attente une fois terminé.
Cette approche présente l'avantage de ne pas nécessiter d'étendre continuellement le bail et de ne pas renvoyer le message tant que l'heure future n'est pas atteinte (ce qui couvre les plantages). Si l'accusé de réception initial réussit, le traitement de type "au plus une fois" est effectué.
Créer des points de contrôle pour les tâches de longue durée
Les files d'attente Spanner peuvent gérer des tâches qui durent de quelques minutes à plusieurs heures, et pas seulement des tâches rapides. Pour ces tâches de longue durée, utilisez l'approche suivante :
- Stocker les métadonnées en externe : utilisez un stockage hors bande pour conserver les détails et l'état de la tâche.
- Créez régulièrement des points de contrôle : pour récupérer après un plantage sans perdre trop de progression, la tâche doit enregistrer régulièrement son état.
- Utilisez le modèle de point de contrôle recommandé : la meilleure façon de créer un point de contrôle consiste à accuser réception de manière atomique (
ACK) du message actuel de la file d'attente et à envoyer un nouveau message dont l'envoi est planifié pour une date ultérieure. Ce nouveau message contient ou pointe vers l'état mis à jour, ce qui empêche la rediffusion immédiate à un autre worker.
Ce modèle réduit le travail en double, même si la création de points de contrôle complets n'est pas possible. Toutefois, dans ce scénario, la tâche redémarre depuis le début après un plantage.
Planifier une tâche pour une date ultérieure spécifique
Pour planifier une tâche à une heure spécifique dans le futur, définissez la colonne DeliverTime lorsque vous insérez le message.
Par exemple, un rappel d'expiration de l'essai :
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');
Gérer les charges utiles de messages volumineuses
Si la charge utile de votre message est volumineuse, utilisez le modèle de stockage hors bande. Stockez la charge utile volumineuse dans une table distincte et ajoutez-y une référence dans le message de la file d'attente.
Par exemple, le traitement des images :
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.
Attendre plusieurs événements avant de continuer
Pour attendre plusieurs événements avant de continuer (comme une opération de jointure), utilisez un tableau pour suivre l'état et une file d'attente pour déclencher les vérifications.
Par exemple, le traitement des commandes nécessitant un inventaire et un paiement :
- Créez une table
OrdersavecInventoryStatusetPaymentStatus. - Lorsque l'inventaire est confirmé, mettez à jour
Orderset envoyez un message àOrderCheckQueue. - Une fois le paiement confirmé, mettez à jour
Orderset envoyez un message àOrderCheckQueue. - Le récepteur de
OrderCheckQueuevérifie la tableOrders. Si les deux états sont confirmés, l'expédition est effectuée et le message est accusé de réception. Sinon, il peut être remis en file d'attente pour une vérification ultérieure ou exécuter une autre logique.
Effectuer une action périodiquement
Pour effectuer une action périodiquement, utilisez le modèle de planification périodique. Le destinataire accuse réception du message et en envoie un nouveau, planifié pour le prochain intervalle.
Par exemple, l'agrégation des données horaires :
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');
Vous pouvez également utiliser les mutations Ack et Send de la bibliothèque cliente. Ces exemples supposent que vous disposez d'un objet Message encapsulant la clé et la charge utile :
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();
});
}
Étapes suivantes
- Découvrez comment utiliser les files d'attente Spanner, y compris les bonnes pratiques et la surveillance.
- Découvrez le traitement de type "exactement une fois" et l'accusé de réception "au plus une fois".
- Configurez le contrôle des accès avec le contrôle des accès précis pour les files d'attente.