Cenários e exemplos de filas do Spanner

Este documento fornece padrões arquitetônicos 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 com vários eventos e criar checkpoints 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 em 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, o registro 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 cinco 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 é reenviada até que o horário futuro chegue (o que inclui falhas). Se o reconhecimento inicial for bem-sucedido, ele vai realizar o processamento "no máximo uma vez".

Criar checkpoints para 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 o estado dela periodicamente.
  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 um 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 colocado em uma nova fila 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 a biblioteca de cliente Ack e as mutações Send. Estes exemplos pressupõem que você tem 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();
  });
}

Agrupar mensagens usando o agrupamento temporal

As filas do Spanner entregam mensagens com latência mínima. No entanto, quando grandes volumes de mensagens chegam continuamente de clientes independentes, o processamento individual de cada mensagem pode criar uma sobrecarga de transação alta. Tentar consultar ou verificar a tabela de filas manualmente para agrupar mensagens pode causar disputa de bloqueio de intervalo, aumento das taxas de interrupção e custos extras de leitura.

Para alcançar o agrupamento de alta capacidade de processamento sem disputa, aplique o padrão de agrupamento temporal. Os remetentes alinham o DeliverTime das mensagens a uma janela de tempo discreta no futuro (por exemplo, arredondando para o limite de 10 segundos mais próximo). Como os remetentes independentes calculam um carimbo de data/hora futuro idêntico, o Spanner tende a agrupar as mensagens da mesma divisão e entregá-las em um único lote se max_batch_size na função com valor de tabela RECEIVE_QUEUE_NAME() permitir ou em vários lotes se max_batch_size for menor que o número de mensagens a serem entregues.

O carimbo de data/hora da entrega pode ser calculado com esta fórmula:

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

Por exemplo, com uma janela de 10 segundos, as mensagens enfileiradas entre 09:05:00 e 09:05:09.999 recebem um DeliverTime de 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)
);

Em seguida, os destinatários extraem essas mensagens com carimbos de data/hora simultâneos em lotes usando 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 você enviar milhões de mensagens ao mesmo tempo, alinhar todas elas ao mesmo segundo exato pode causar picos repentinos de processamento. Para distribuir o trabalho de maneira uniforme e ainda agrupar mensagens por entidade, adicione um deslocamento ao cálculo com base em um identificador exclusivo (como TIMESTAMP_ADD(deliver_time, INTERVAL MOD(ABS(FARM_FINGERPRINT(CAST(UserId AS STRING))), 10) SECOND)).

Desvincular a modificação de dados do processamento (padrão de flag suja)

Em aplicativos transacionais de alta taxa de transferência, executar recálculos complexos, indexação de pesquisa ou invalidação de cache diretamente em transações voltadas ao usuário pode diminuir a velocidade da experiência do usuário, aumentando a latência e causando disputa de bloqueio em linhas compartilhadas.

O padrão de flag suja desacopla as modificações de dados do processamento assíncrono. Quando uma transação modifica uma tabela, ela grava uma mensagem leve de "bit sujo" em uma fila na mesma transação. Um worker em segundo plano consome a mensagem e realiza o processamento caro de forma assíncrona.

Por exemplo, quando um cliente faz uma mudança no perfil ou nas configurações:

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

O receptor em segundo plano para UserDirtyQueue recebe o UserId, lê a linha de perfil atualizada fora do caminho crítico do usuário e recalcula o índice de pesquisa ou atualiza os caches externos. O agrupamento em lote também pode ser aplicado a esse padrão se várias atualizações do mesmo usuário forem enviadas à fila em um curto período. Nesse caso, a TVF do receptor pode especificar um max_batch_size maior que 1 para receber várias mensagens do mesmo lote.

Monitorar a integridade do worker e detectar tempos limite

É possível usar mensagens de fila programadas para criar um sistema de pulsação tolerante a falhas e de monitoramento de integridade para frotas de nós de trabalho ou instâncias de microsserviços.

Para implementar a verificação de integridade:

  1. Registrar worker na inicialização:quando um worker é inicializado, ele insere uma mensagem de pulsação em uma fila de verificação de integridade com um DeliverTime futuro definido para o prazo de falha (por exemplo, 60 segundos).
  2. Enviar pulsações periódicas:enquanto estiver funcionando corretamente, o worker vai atualizar periodicamente (por exemplo, a cada 10 segundos) a mensagem de pulsação, avançando o DeliverTime mais 60 segundos no futuro.
  3. Detectar falhas:se o worker falhar ou perder a conectividade de rede, as atualizações de pulsação serão interrompidas. Após 60 segundos, o carimbo de data/hora de entrega é atualizado (DeliverTime <= CURRENT_TIMESTAMP()), e o Spanner entrega a mensagem a um receptor de alertas, que inicia o failover ou a reatribuição de tarefas.

Importante:as filas do Spanner não são compatíveis com instruções de DML UPDATE. Portanto, para atualizar o carimbo de data/hora de pulsação, exclua a mensagem atual e insira uma substituição com o novo DeliverTime em uma única transação ou aplique mutações Ack e Send da biblioteca de cliente. Verifique se a chave primária da fila é apenas WorkerId (em vez de (WorkerId, MessageId)) para que haja apenas uma mensagem de pulsação por worker a qualquer 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)
);

Ao usar mutações da biblioteca de cliente, verifique se Ack precede Send na fatia de mutação, como neste exemplo em 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 um worker é desligado normalmente, ele exclui explicitamente a mensagem de pulsação para que nenhum alerta falso seja acionado:

DELETE FROM WorkerHealthQueue WHERE WorkerId = 'worker-node-42' ASSERT_ROWS_MODIFIED 1;

Processar vários tipos de tarefas em uma única fila

As instâncias do Spanner têm limites no número total de filas. Criar uma fila separada para cada pequena operação assíncrona pode atingir esse limite rapidamente e exige o gerenciamento de muitas consultas de receptor simultâneas.

Para consolidar operações, combine diferentes tipos de tarefas em uma única fila, chamada de fila polimórfica. Há duas estratégias usadas para criar uma fila polimórfica.

Estratégia 1: incluir uma coluna de tipo na chave primária

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

O receptor inspeciona TaskType e envia o payload para o manipulador correspondente.

Estratégia 2: estrutura de payload polimórfica

Como alternativa, use um payload JSON que contenha um campo de ação ou discriminador de tipo:

{
  "action": "SYNC_INVENTORY",
  "data": { "item_id": 987, "delta": -1 }
}

Implementar atrasos personalizados para novas tentativas

As filas do Spanner tentam automaticamente de novo as mensagens com falha ou não confirmadas com espera exponencial integrada. No entanto, em cenários em que uma mensagem falha por um motivo conhecido com uma duração conhecida ou um limite de taxa externa (como uma resposta HTTP 429 especificando um cabeçalho Retry-After), confiar no backoff automático pode causar tentativas prematuras que desperdiçam recursos de CPU.

Para implementar um atraso de nova tentativa personalizado:

  1. Detecte a falha transitória específica no seu processador de mensagens.
  2. Confirme a mensagem atual para satisfazer a tentativa de entrega atual.
  3. Na mesma transação, envie uma mensagem de substituição com um DeliverTime explícito definido para o horário de nova tentativa escolhido (no exemplo a seguir, 5 minutos depois).

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

Essa abordagem permite que o aplicativo gerencie com precisão as programações de espera e evite saturar APIs externas durante períodos de recuperação downstream.

A seguir