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:
- O destinatário recebe uma mensagem de
RECEIVE_ReportQueue(). - Inicie a geração do relatório.
- A cada 5 segundos, chame
SELECT * FROM RENEWLEASE_ReportQueue([leaseToken])em uma linha de execução ou rotina separada. - 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:
- Reconheça (
DELETEouACK) 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. - 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:
- Armazenar metadados externamente:use o armazenamento fora da banda para manter os detalhes e o estado da tarefa.
- Crie checkpoints regularmente:para se recuperar de falhas sem perder muito progresso, a tarefa precisa salvar periodicamente o estado dela.
- 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:
- Crie uma tabela
OrderscomInventoryStatusePaymentStatus. - Quando o inventário for confirmado, atualize
Orderse envie uma mensagem paraOrderCheckQueue. - Quando o pagamento for confirmado, atualize
Orderse envie uma mensagem paraOrderCheckQueue. - O receptor de
OrderCheckQueueverifica a tabelaOrders. 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
- Saiba como usar filas do Spanner, incluindo práticas recomendadas e monitoramento.
- Saiba mais sobre o processamento "exatamente uma vez" e o reconhecimento "no máximo uma vez".
- Configure o controle de acesso com o controle de acesso detalhado para filas.