Scenari ed esempi di code Spanner

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:

  1. Il destinatario riceve un messaggio da RECEIVE_ReportQueue().
  2. Avvia la generazione del report.
  3. Ogni 5 secondi, chiama SELECT * FROM RENEWLEASE_ReportQueue([leaseToken]) in un thread o una routine separati.
  4. 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:

  1. Riconosci (DELETE o ACK) 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.
  2. 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:

  1. Archivia metadati esternamente:utilizza l'archiviazione out-of-band per conservare i dettagli e lo stato dell'attività.
  2. 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.
  3. 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:

  1. Crea una tabella Orders con InventoryStatus e PaymentStatus.
  2. Una volta confermato l'inventario, aggiorna Orders e invia un messaggio a OrderCheckQueue.
  3. Una volta confermato il pagamento, aggiorna Orders e invia un messaggio a OrderCheckQueue.
  4. Il destinatario di OrderCheckQueue controlla la tabella Orders. 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:

$$ \text{DeliverTime} = \text{RoundDown}(\text{CurrentTime}, \text{FixedDelay}) + \text{FixedDelay} $$

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à:

  1. Registra worker all'avvio:quando un worker viene inizializzato, inserisce un messaggio heartbeat in una coda di controllo dell'integrità con un DeliverTime futuro impostato alla scadenza del guasto (ad esempio, 60 secondi).
  2. Invia battiti periodici: quando è integro, il worker aggiorna periodicamente (ad esempio ogni 10 secondi) il messaggio di battito incrementando il valore di DeliverTime di altri 60 secondi nel futuro.
  3. 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:

  1. Rileva l'errore temporaneo specifico nel processore di messaggi.
  2. Conferma il messaggio corrente per soddisfare il tentativo di consegna corrente.
  3. Nella stessa transazione, invia un messaggio di sostituzione con un DeliverTime esplicito 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