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 l'archiviazione out-of-band, coordinare flussi di lavoro multi-evento ed eseguire il checkpoint o estendere i lease per i 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)
);
Una volta eseguito 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 di lunga durata che richiede l'elaborazione al massimo una volta o un periodo di lease lungo, procedi nel seguente modo:
- Riconosci (
DELETEoACK) il messaggio in coda corrente all'arrivo. Nella stessa transazione, metti di nuovo in coda un nuovo messaggio di 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 nuovamente fino al momento futuro (che copre gli arresti anomali). Se l'acknowledgement iniziale va a buon fine, viene eseguita l'elaborazione al massimo una volta.
Creare checkpoint per 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 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 in caso di 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 il checkpoint completo non è possibile, anche se in questo scenario l'attività viene riavviata dall'inizio dopo un arresto anomalo.
Pianificare il lavoro per un momento specifico nel futuro
Per programmare il lavoro per un momento specifico in futuro, imposta la colonna DeliverTime
quando inserisci il messaggio.
Ad esempio, un promemoria di 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 unione), 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 conferma la ricezione del 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();
});
}
Raggruppa i messaggi in batch utilizzando il raggruppamento temporale
Le code Spanner inviano i messaggi con una latenza minima. Tuttavia, quando arrivano continuamente volumi elevati di messaggi da client indipendenti, l'elaborazione di ogni messaggio singolarmente può creare un sovraccarico elevato delle transazioni. Il tentativo di eseguire query o scansioni manuali della tabella della coda per raggruppare i messaggi può introdurre contese di blocco dell'intervallo, tassi di interruzione elevati e costi di lettura aggiuntivi.
Per ottenere il batching a velocità effettiva elevata senza contesa, applica il pattern di batching
temporale. I mittenti allineano il DeliverTime dei messaggi a una finestra temporale discreta in futuro (ad esempio, arrotondando al limite di 10 secondi più vicino). Poiché i mittenti indipendenti calcolano un timestamp futuro identico,
Spanner tende a raggruppare i messaggi della stessa suddivisione
e li invia in un unico batch se max_batch_size nella
funzione con valori di tabella RECEIVE_QUEUE_NAME()
lo consente o in più batch se max_batch_size è inferiore al numero di
messaggi da inviare.
Il timestamp di pubblicazione può essere calcolato con questa formula:
Ad esempio, con una finestra di 10 secondi, i messaggi accodati tra le 09:05:00 e le
09:05:09.999 ricevono tutti un DeliverTime di 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)
);
I destinatari estraggono quindi questi messaggi sincronizzati in batch utilizzando max_batch_size:
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');
Se invii milioni di messaggi contemporaneamente, l'allineamento di tutti i messaggi allo stesso secondo esatto può causare picchi di elaborazione improvvisi. Per distribuire il lavoro
in modo uniforme e raggruppare comunque i messaggi per entità, aggiungi un offset al calcolo
in base a un identificatore univoco (ad esempio TIMESTAMP_ADD(deliver_time, INTERVAL MOD(ABS(FARM_FINGERPRINT(CAST(UserId AS STRING))), 10) SECOND)).
Disaccoppiare la modifica dei dati dall'elaborazione (pattern dirty flag)
Nelle applicazioni transazionali a velocità effettiva elevata, l'esecuzione di ricalcoli complessi, l'indicizzazione della ricerca o l'invalidazione della cache direttamente all'interno delle transazioni rivolte agli utenti può rallentare l'esperienza utente aumentando la latenza e causando contese di blocco sulle righe condivise.
Il pattern del flag sporco disaccoppia le modifiche ai dati dall'elaborazione asincrona. Quando una transazione modifica una tabella, scrive un messaggio "dirty bit" leggero in una coda all'interno della stessa transazione. Un worker in background utilizza successivamente il messaggio ed esegue l'elaborazione costosa in modo asincrono.
Ad esempio, quando un cliente modifica un profilo o le impostazioni:
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));
Il ricevitore in background per UserDirtyQueue riceve UserId, legge la
riga del profilo aggiornata al di fuori del percorso critico dell'utente e ricalcola l'indice di ricerca
o aggiorna le cache esterne. Il batch può essere applicato anche a questo pattern se più aggiornamenti dello stesso utente vengono inviati alla coda in un breve periodo di tempo. In questo caso, il TVF ricevente può specificare un
max_batch_size maggiore di 1 per ricevere più messaggi dallo stesso
batch.
Monitora l'integrità dei worker e rileva i timeout
Puoi utilizzare i messaggi in coda pianificati per creare un sistema di heartbeat e monitoraggio dell'integrità a tolleranza di errore per parchi risorse di nodi worker o istanze di microservizi.
Per implementare il controllo di integrità:
- Registra worker all'avvio:quando un worker viene inizializzato, inserisce un
messaggio heartbeat in una coda di controllo dell'integrità con un
DeliverTimefuturo impostato alla scadenza del guasto (ad esempio, 60 secondi). - Invia battiti periodici: quando è integro, il worker aggiorna periodicamente (ad esempio ogni 10 secondi) il messaggio di battito incrementando il valore di
DeliverTimedi altri 60 secondi nel futuro. - Rileva errori:se il worker si arresta in modo anomalo o perde la connettività di rete,
gli aggiornamenti heartbeat si interrompono. Dopo 60 secondi, il timestamp di consegna matura
(
DeliverTime <= CURRENT_TIMESTAMP()) e Spanner consegna il messaggio a un destinatario di avvisi, che avvia il failover o la riassegnazione dell'attività.
Importante:le code Spanner non supportano le istruzioni DML UPDATE. Pertanto, per aggiornare il timestamp heartbeat, devi eliminare il messaggio esistente e inserirne uno sostitutivo con il nuovo DeliverTime all'interno di una singola transazione oppure applicare le mutazioni delle librerie client Ack e Send. Assicurati
che la chiave primaria della coda sia solo WorkerId (anziché
(WorkerId, MessageId)) in modo che esista un solo messaggio heartbeat per worker in
un determinato momento.
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)
);
Quando utilizzi le mutazioni della libreria client, assicurati che Ack preceda Send nella
sezione delle mutazioni, come in questo esempio Go:
// 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
}
Quando un worker si arresta in modo controllato, elimina esplicitamente il messaggio heartbeat in modo che non venga attivato alcun falso avviso:
DELETE FROM WorkerHealthQueue WHERE WorkerId = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;
Gestire più tipi di attività in un'unica coda
Le istanze Spanner hanno limiti al numero totale di code. La creazione di una coda separata per ogni piccola operazione asincrona può raggiungere rapidamente questo limite e richiede la gestione di molte query del destinatario simultanee.
Per consolidare le operazioni, combina diversi tipi di attività in un'unica coda, chiamata coda polimorfica. Esistono due strategie utilizzate per creare una coda polimorfica.
Strategia 1: includi una colonna di tipo nella chiave primaria
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));
Il ricevitore esamina TaskType e invia il payload al gestore corrispondente.
Strategia 2: struttura del payload polimorfico
In alternativa, utilizza un payload JSON contenente un campo discriminatore di azione o tipo:
{
"action": "SYNC_INVENTORY",
"data": { "item_id": 987, "delta": -1 }
}
Implementare ritardi personalizzati per i nuovi tentativi
Spanner Queues ritenta automaticamente l'invio dei messaggi non riusciti o non riconosciuti
con il backoff esponenziale integrato. Tuttavia, negli scenari in cui un messaggio
non viene recapitato per un motivo noto con una durata nota o
un limite di frequenza esterno (ad esempio una risposta HTTP 429 che specifica un'intestazione Retry-After), fare affidamento sul backoff automatico può causare tentativi di ripetizione prematuri che
sprecano risorse della CPU.
Per implementare un ritardo personalizzato per i nuovi tentativi:
- Rileva l'errore temporaneo specifico nel processore di messaggi.
- Conferma il messaggio corrente per soddisfare il tentativo di consegna corrente.
- Nella stessa transazione, invia un messaggio di sostituzione con un
DeliverTimeesplicito impostato sull'ora di nuovo tentativo futura scelta (nell'esempio seguente, 5 minuti dopo).
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)
);
Questo approccio consente all'applicazione di gestire con precisione le pianificazioni di backoff ed evitare di saturare le API esterne durante i periodi di recupero downstream.
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.