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

  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 a lunga esecuzione che richiede l'elaborazione al massimo una volta o un periodo di lease lungo, procedi nel seguente modo:

  1. Al tuo arrivo, conferma (DELETE o ACK) 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.
  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 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:

  1. Archivia i 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 dagli 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 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:

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