Cenários e exemplos de filas do Spanner

Este documento fornece padrões de arquitetura e exemplos de código para cenários comuns de mensagens usando filas do Spanner. É possível usar esses padrões para acionar trabalhos assíncronos após o commit de transações, programar tarefas atrasadas ou recorrentes, gerenciar grandes payloads de mensagens com armazenamento fora da banda, coordenar fluxos de trabalho de vários eventos e criar pontos de verificação ou estender concessões para jobs em segundo plano de longa duração.

Processamento único e confirmação de recebimento no máximo uma vez

As várias considerações e soluções para o processamento único e o reconhecimento de no máximo uma vez são descritas com mais detalhes na página Processamento único e reconhecimento de no máximo uma vez.

Realizar o trabalho depois que uma transação é confirmada

Para realizar o trabalho depois que uma transação é confirmada, envie uma mensagem para a fila na mesma transação.

Por exemplo, a inscrição de um novo usuário aciona um e-mail de boas-vindas:

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

Depois que a transação é confirmada, o receptor de UserTasks transmite a mensagem, envia o e-mail e confirma a mensagem:

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

Processar trabalho de longa duração

Se você tiver um trabalho que possa levar mais tempo do que o período de concessão padrão (mais de 10 segundos), chame SELECT * FROM RENEWLEASE_QUEUE_NAME() periodicamente.

Por exemplo, gerar um relatório:

  1. O destinatário recebe uma mensagem de RECEIVE_ReportQueue().
  2. Inicie a geração do relatório.
  3. A cada 5 segundos, chame SELECT * FROM RENEWLEASE_ReportQueue([leaseToken]) em uma linha de execução ou rotina separada.
  4. Quando terminar, confirme a mensagem e armazene o relatório.

Como alternativa, se você tiver um trabalho de longa duração que exija o processamento no máximo uma vez ou um tempo de concessão longo, faça o seguinte:

  1. Reconheça (DELETE ou ACK) a mensagem atual da fila ao chegar. Na mesma transação, coloque novamente na fila uma nova mensagem com um carimbo de data/hora de entrega no futuro, além do tempo necessário para o processamento.
  2. Continue o processamento e confirme a mensagem recém-enfileirada quando concluído.

As vantagens dessa abordagem são que não é necessário estender continuamente o aluguel, e a mensagem não é entregue novamente até que o horário futuro chegue (o que cobre falhas). Se o reconhecimento inicial for bem-sucedido, ele vai alcançar o processamento "no máximo uma vez".

Criar checkpoints de trabalhos de longa duração

As filas do Spanner podem gerenciar tarefas que duram de minutos a horas, não apenas jobs rápidos. Para essas tarefas de longa duração, use a seguinte abordagem:

  1. Armazenar metadados externamente:use o armazenamento fora da banda para manter os detalhes e o estado da tarefa.
  2. Crie checkpoints regularmente:para se recuperar de falhas sem perder muito progresso, a tarefa precisa salvar periodicamente o estado dela.
  3. Use o padrão de checkpoint recomendado:a melhor maneira de fazer isso é reconhecer atomicamente (ACK) a mensagem atual da fila e enviar uma nova mensagem programada para entrega futura. Essa nova mensagem contém ou aponta para o estado atualizado, o que impede a nova entrega imediata a outro worker.

Esse padrão reduz o trabalho duplicado, mesmo que não seja possível fazer um checkpoint completo, embora a tarefa seja reiniciada do início após uma falha nesse cenário.

Programar trabalho para um horário específico no futuro

Para programar o trabalho para um horário específico no futuro, defina a coluna DeliverTime ao inserir a mensagem.

Por exemplo, um lembrete de expiração do teste:

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');

Processar payloads de mensagens grandes

Se o payload da mensagem for grande, use o padrão de armazenamento fora da banda. Armazene o payload grande em uma tabela separada e coloque uma referência a ele na mensagem da fila.

Por exemplo, processamento de imagens:

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.

Aguardar vários eventos antes de continuar

Para aguardar vários eventos antes de continuar (como uma operação de junção), use uma tabela para rastrear o estado e uma fila para acionar verificações.

Por exemplo, o atendimento de pedidos que exigem inventário e pagamento:

  1. Crie uma tabela Orders com InventoryStatus e PaymentStatus.
  2. Quando o inventário for confirmado, atualize Orders e envie uma mensagem para OrderCheckQueue.
  3. Quando o pagamento for confirmado, atualize Orders e envie uma mensagem para OrderCheckQueue.
  4. O receptor de OrderCheckQueue verifica a tabela Orders. Se os dois status forem confirmados, o processo de envio vai continuar e a mensagem será reconhecida. Caso contrário, ele poderá ser enfileirado novamente para uma verificação posterior ou executar outra lógica.

Realizar uma ação periodicamente

Para realizar uma ação periodicamente, use o padrão de programação periódica. O receptor confirma a mensagem e envia uma nova programada para o próximo intervalo.

Por exemplo, agregação de dados por hora:

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');

Como alternativa, use as mutações Ack e Send da biblioteca de cliente. Estes exemplos pressupõem que você tenha um objeto Message encapsulando a chave e o 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();
  });
}

A seguir