Questo documento fornisce pattern architetturali ed esempi di codice per scenari di messaggistica comuni utilizzando le code Spanner. Puoi utilizzare questi pattern per attivare il lavoro asincrono dopo il commit delle transazioni, pianificare attività ritardate o ricorrenti, gestire payload di messaggi di grandi dimensioni con spazio di archiviazione out-of-band, coordinare flussi di lavoro multi-evento ed eseguire checkpoint o estendere i lease per job in background a lunga esecuzione.
Elaborazione "exactly-once" e riconoscimento "at-most-once"
Le varie considerazioni e soluzioni per l'elaborazione "exactly-once" e il riconoscimento "at-most-once" sono descritte in modo più dettagliato nella pagina Elaborazione "exactly-once" e riconoscimento "at-most-once".
Esegui il lavoro dopo il commit di una transazione
Per eseguire il lavoro dopo il commit di una transazione, invia un messaggio alla coda all'interno della stessa transazione.
Ad esempio, la registrazione di un nuovo utente attiva un'email di benvenuto:
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)
);
Dopo il commit della transazione, il destinatario di UserTasks trasmette in streaming il
messaggio, invia l'email e conferma la ricezione del messaggio:
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';
Gestire le attività di lunga durata
Se hai un lavoro che potrebbe richiedere più tempo del lease predefinito (più di 10
secondi), chiama periodicamente SELECT * FROM RENEWLEASE_QUEUE_NAME().
Ad esempio, la generazione di un report:
- Il destinatario riceve un messaggio da
RECEIVE_ReportQueue(). - Avvia la generazione del report.
- Ogni 5 secondi, chiama
SELECT * FROM RENEWLEASE_ReportQueue([leaseToken])in un thread o una routine separati. - Al termine, conferma il messaggio e archivia il report.
In alternativa, se hai un lavoro a lunga esecuzione che richiede l'elaborazione al massimo una volta o un periodo di lease lungo, procedi nel seguente modo:
- Al tuo arrivo, conferma (
DELETEoACK) il messaggio in coda corrente. Nella stessa transazione, rimetti in coda un nuovo messaggio della coda con un timestamp di consegna futuro, oltre il tempo necessario per l'elaborazione. - Procedi con l'elaborazione e conferma il messaggio appena messo in coda al termine dell'operazione.
I vantaggi di questo approccio sono che non è necessario estendere continuamente il lease e il messaggio non viene inviato di nuovo fino al momento futuro (che copre gli arresti anomali). Se l'acknowledgement iniziale ha esito positivo, viene eseguita l'elaborazione al massimo una volta.
Controllare le attività di lunga durata
Le code Spanner possono gestire attività che durano da minuti a ore, non solo job rapidi. Per queste attività di lunga durata, utilizza il seguente approccio:
- Archivia i metadati esternamente:utilizza l'archiviazione out-of-band per conservare i dettagli e lo stato dell'attività.
- Esegui regolarmente il checkpoint:per eseguire il ripristino dagli arresti anomali senza perdere molti progressi, l'attività deve salvare periodicamente il proprio stato.
- Utilizza il pattern di checkpointing consigliato:il modo migliore per eseguire il checkpointing è
riconoscere in modo atomico (
ACK) il messaggio corrente della coda e inviare un nuovo messaggio pianificato per la consegna futura. Questo nuovo messaggio contiene o punta allo stato aggiornato, il che impedisce la nuova consegna immediata a un altro worker.
Questo pattern riduce il lavoro duplicato anche se non è possibile il checkpointing completo, anche se l'attività viene riavviata dall'inizio dopo un arresto anomalo in questo scenario.
Pianificare il lavoro per un momento specifico nel futuro
Per programmare un'attività per un orario specifico in futuro, imposta la colonna DeliverTime
quando inserisci il messaggio.
Ad esempio, un promemoria relativo alla scadenza della prova:
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');
Gestire payload di messaggi di grandi dimensioni
Se il payload del messaggio è grande, utilizza il pattern di archiviazione out-of-band. Archivia il payload di grandi dimensioni in una tabella separata e inserisci un riferimento nel messaggio della coda.
Ad esempio, l'elaborazione delle immagini:
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.
Attendi più eventi prima di procedere
Per attendere più eventi prima di procedere (ad esempio un'operazione di join), utilizza una tabella per monitorare lo stato e una coda per attivare i controlli.
Ad esempio, l'evasione dell'ordine che richiede inventario e pagamento:
- Crea una tabella
OrdersconInventoryStatusePaymentStatus. - Una volta confermato l'inventario, aggiorna
Orderse invia un messaggio aOrderCheckQueue. - Una volta confermato il pagamento, aggiorna
Orderse invia un messaggio aOrderCheckQueue. - Il destinatario di
OrderCheckQueuecontrolla la tabellaOrders. Se entrambi gli stati sono confermati, la spedizione procede e il messaggio viene riconosciuto. In caso contrario, potrebbe essere rimesso in coda per un controllo successivo o eseguire un'altra logica.
Eseguire un'azione periodicamente
Per eseguire un'azione periodicamente, utilizza il pattern di pianificazione periodica. Il ricevitore riconosce il messaggio e ne invia uno nuovo programmato per l'intervallo successivo.
Ad esempio, l'aggregazione dei dati orari:
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');
In alternativa, utilizza le mutazioni Ack e Send della libreria client. Questi esempi presuppongono che tu disponga di un oggetto Message che incapsula la chiave e il payload:
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();
});
}
Passaggi successivi
- Scopri come utilizzare le code Spanner, incluse le best practice e il monitoraggio.
- Scopri di più sull'elaborazione "exactly-once" e sull'acknowledgement "at-most-once".
- Configura il controllo dell'accesso con il controllo dell'controllo dell'accesso granulare per le code.