Este documento descreve como usar filas do Spanner. Ele explica como criar uma fila, enviar e receber mensagens, estender concessões de mensagens e confirmar mensagens. Ele também inclui práticas recomendadas, informações sobre filas de monitoramento e orientações para solução de problemas.
Crie uma fila
Para criar uma fila, use a instrução CREATE QUEUE.
GoogleSQL
-- Example table for interleaving
CREATE TABLE Users (
UserId INT64 NOT NULL,
UserName STRING(MAX)
) PRIMARY KEY (UserId);
-- Queue for processing user-related tasks
-- NOTE: Queues do not have to be interleaved, but for locality and
-- pre-warming purposes, it's recommended if you are inserting into a table
-- and a queue simultaneously.
CREATE QUEUE UserTasks (
UserId INT64 NOT NULL,
MessageId STRING(36) NOT NULL, -- UUID recommended
Payload BYTES(MAX) NOT NULL -- Also: Proto, JSON, String are possible.
) PRIMARY KEY (UserId, MessageId),
INTERLEAVE IN PARENT Users ON DELETE CASCADE;
PostgreSQL
-- Example table for interleaving
CREATE TABLE users (
userid bigint NOT NULL,
username varchar,
PRIMARY KEY (userid)
);
-- Queue for processing user-related tasks
-- NOTE: Queues do not have to be interleaved, but for locality and
-- pre-warming purposes, it's recommended if you are inserting into a table
-- and a queue simultaneously.
CREATE QUEUE usertasks (
userid bigint NOT NULL,
messageid varchar(36) NOT NULL, -- UUID recommended
payload bytea NOT NULL, -- Also: text, varchar, jsonb are possible.
PRIMARY KEY (userid, messageid)
) INTERLEAVE IN PARENT users ON DELETE CASCADE;
A única coluna que não precisa ser criada explicitamente na instrução CREATE
QUEUE é chamada de DeliverTime no GoogleSQL e deliver_time no PostgreSQL. Elas são criadas automaticamente pelo Spanner.
As filas aceitam políticas de time to live (TTL), que podem ajudar a gerenciar o backlog de mensagens antigas e não confirmadas.
Enviar uma mensagem
Para enviar uma mensagem a uma fila, use a instrução
INSERT DML:
GoogleSQL
-- Send message immediately
INSERT INTO Users (UserId) VALUES (123);
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (123, 'some-unique-id-1', b'Your task payload here', CURRENT_TIMESTAMP());
-- Schedule a message delivery
INSERT INTO UserTasks (UserId, MessageId, Payload, DeliverTime)
VALUES (123, 'some-unique-id-2', b'Scheduled task', TIMESTAMP_ADD(CURRENT_TIMESTAMP(), INTERVAL 1 HOUR));
PostgreSQL
-- Send message immediately
INSERT INTO users (userid) VALUES (123);
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (123, 'some-unique-id-1', CAST('Your task payload here' AS bytea), CURRENT_TIMESTAMP);
-- Schedule a message delivery
INSERT INTO usertasks (userid, messageid, payload, deliver_time)
VALUES (123, 'some-unique-id-2', CAST('Scheduled task' AS bytea), CURRENT_TIMESTAMP + INTERVAL '1 HOUR');
Ou use as mutações da biblioteca de cliente Send para inserir uma mensagem:
Go
m := spanner.Send("UserTasks", spanner.Key{int64(123), "some-unique-id-1"}, []byte("Your task payload here"), spanner.WithDeliveryTime(futureTime))
_, err := client.Apply(ctx, []*spanner.Mutation{m})
Java
dbClient.write(
Collections.singletonList(
Mutation.newSendBuilder("UserTasks")
.setKey(Key.of(123L, "some-unique-id-1"))
.setPayload(Value.bytes(ByteArray.copyFrom("message3")))
.setDeliveryTime(futureTime)
.build()));
Receber mensagens
Use a função com valor de tabela (TVF) RECEIVE_QUEUE_NAME() com ExecuteStreamingSQL para receber mensagens. Esta é uma chamada de longa duração. É necessário executar uma dessas chamadas por worker, por fila, de maneira repetida. Execute a consulta usando uma leitura consistente porque o Spanner rejeita leituras obsoletas.
GoogleSQL
-- SQL query to stream messages
SELECT
UserId,
MessageId,
Payload,
DeliverTime,
SpannerLeaseExpirationTimestamp,
SpannerLeaseToken
FROM RECEIVE_UserTasks(max_duration=>'20m');
PostgreSQL
-- SQL query to stream messages
SELECT
userid,
messageid,
payload,
deliver_time,
spanner_lease_expiration_timestamp,
spanner_lease_token
FROM spanner.receive_usertasks(NULL, NULL, '20m');
O código do cliente precisa iterar os resultados usando uma consulta de streaming. Cada linha retornada é uma mensagem.
Receber mensagens em lote
Para aumentar a capacidade de processamento de várias mensagens juntas, é possível receber
mensagens em lotes especificando o argumento max_batch_size:
GoogleSQL
-- SQL query to stream messages
SELECT
UserId,
MessageId,
Payload,
DeliverTime,
SpannerLeaseExpirationTimestamp,
SpannerLeaseToken,
SpannerLastBatchMessage -- Special boolean column returns TRUE if the
-- last message is in a batch.
FROM RECEIVE_UserTasks(max_duration=>'20m', max_batch_size=>20);
PostgreSQL
-- SQL query to stream messages
SELECT
userid,
messageid,
payload,
deliver_time,
spanner_lease_expiration_timestamp,
spanner_lease_token,
spanner_last_batch_message -- Special boolean column returns TRUE if the
-- last message is in a batch.
FROM spanner.receive_usertasks(20, NULL, '20m');
Se preferir, use a biblioteca de cliente. Este exemplo em Go demonstra como transmitir mensagens de uma fila, verificar o vencimento do aluguel e confirmar mensagens de forma assíncrona:
// import "cloud.google.com/go/spanner"
// import "google.golang.org/api/iterator"
stmt := spanner.Statement{SQL: "SELECT * FROM RECEIVE_UserTasks(max_duration=>'20m')"}
iter := client.Single().Query(ctx, stmt)
defer iter.Stop()
for {
row, err := iter.Next()
if err == iterator.Done {
break // Or potentially restart the query
}
if err != nil {
// Handle error
return err
}
var userId int64
var messageId string
var payload []byte
var deliverTime time.Time
var leaseExpiration time.Time
var leaseToken string
var lastBatchMessage bool
if err := row.Columns(&userId, &messageId, &payload, &deliverTime, &leaseExpiration, &leaseToken, &lastBatchMessage); err != nil {
// Handle column parsing error
return err
}
if time.Now().After(leaseExpiration) {
log.Printf("Lease expired for message %s, skipping", messageId)
continue
}
// Process and acknowledge the message asynchronously so that we don't
// block receiving subsequent messages.
go func(userId int64, messageId string, payload []byte) {
// Process the message (payload)
// ... potentially long-running work ...
// Need to extend the lease if processing is long
// Acknowledge the message upon success
_, err := client.Apply(ctx, []*spanner.Mutation{
spanner.Ack("UserTasks", spanner.Key{userId, messageId}),
})
if err != nil {
// Handle ack error
}
}(userId, messageId, payload)
}
Estender o período de concessão da mensagem
Se o processamento de uma mensagem levar mais tempo do que a concessão inicial (10 segundos), use a seguinte sintaxe:
GoogleSQL
-- Extends the lease and returns new expiration time lease_token.
-- The lease will be extended by 10s (not configurable).
SELECT * FROM RENEWLEASE_UserTasks(lease_tokens => ['token1', ..., 'tokenN'])
-- Returns rows of tokens and whether they were successfully extended
SpannerOldLeaseToken SpannerNewLeaseToken SpannerLeaseExpirationTimestamp
<old_token1> <new_token1> "2025-09-27T12:10:00.0Z"
'token2' <NULL> "2025-09-27T12:09:51.0Z"
...
'tokenN' 'tokenN_new' "2025-09-27T12:10:01.0Z"
PostgreSQL
-- Extends the lease and returns new expiration time lease_token.
-- The lease will be extended by 10s (not configurable).
SELECT * FROM spanner.renewlease_usertasks(lease_tokens => ARRAY['token1', ..., 'tokenN'])
-- Returns rows of tokens and whether they were successfully extended
spanner_old_lease_token spanner_new_lease_token spanner_lease_expiration_timestamp
<old_token1> <new_token1> "2025-09-27T12:10:00.0Z"
'token2' <NULL> "2025-09-27T12:09:51.0Z"
...
'tokenN' 'tokenN_new' "2025-09-27T12:10:01.0Z"
Os tokens de concessão são retornados de acordo com a seguinte lógica:
- Tokens de concessão não analisáveis não retornam uma linha.
- Tokens de concessão já expirados não retornam uma linha.
- Tokens de concessão não renováveis não retornam uma linha com um
SpannerNewLeaseTokende NULL. Isso pode acontecer se a mensagem já tiver sido confirmada, mas o token de concessão não tiver expirado.
Se preferir, use a biblioteca de cliente. Este exemplo em Go demonstra como estender uma concessão de mensagem:
// Inside message processing loop...
// Before leaseExpiration, for example, in a separate goroutine or timed check
// leaseToken is from the SELECT query
extendStmt := spanner.Statement{
SQL: "SELECT * FROM RENEWLEASE_UserTasks(lease_tokens => [@token1])",
Params: map[string]interface{}{
"token1": leaseToken,
},
}
_, err := client.Single().Query(ctx, extendStmt).Next() // Simplified call
if err != nil {
log.Printf("Failed to extend lease for %s: %v", messageId, err)
// Processing should probably stop as redelivery is likely
} else {
// New lease expiration is typically approximately 10s from now
log.Printf("Lease extended for %s", messageId)
// Update local leaseExpiration time if needed
}
Confirmar uma mensagem
Use o DML DELETE
para confirmar uma mensagem. Isso precisa ser transacional com qualquer outra gravação
relacionada ao processamento de mensagens.
GoogleSQL
-- DML for acknowledging
DELETE FROM UserTasks WHERE UserId = @userId AND MessageId = @messageId ASSERT_ROWS_MODIFIED 1;
PostgreSQL
-- DML for acknowledging
DELETE FROM usertasks WHERE userid = $1 AND messageid = $2 ASSERT_ROWS_MODIFIED 1;
Como alternativa, use a mutação Ack da biblioteca de cliente para
confirmar uma mensagem:
Go
_, err := client.Apply(ctx, []*spanner.Mutation{
spanner.Ack("UserTasks", spanner.Key{1}),
})
Java
dbClient.write(
Collections.singletonList(
Mutation.newAckBuilder("UserTasks")
.setKey(Key.of(2L))
.build()));
Práticas recomendadas
Confira a seguir as práticas recomendadas para usar filas do Spanner:
- Payloads pequenos:mantenha os payloads de mensagens da fila pequenos, com menos de 4 KB. Use o padrão de armazenamento fora da banda para dados maiores.
- Gerenciamento de concessão:estenda as concessões de tarefas que podem exceder a concessão de mensagem padrão de 10 segundos. Se você não estender os contratos de mensagens, poderá haver novas entregas e um possível processamento duplo.
- Tratamento de erros:o Spanner enfileira mensagens de nova tentativa que não são processadas com espera exponencial na primeira hora, e as mensagens mais antigas são repetidas uma vez por hora. Considere mover as mensagens que falham permanentemente para uma fila separada.
- Monitoramento:monitore as profundidades da fila e a idade da mensagem mais antiga não confirmada para garantir que os receptores não estejam limitando a entrada do pipeline.
- Idempotência:projete seus processadores de mensagens para tolerar operações repetidas, porque a entrega do tipo "pelo menos uma vez" significa que duplicatas ocasionais são possíveis. Consulte a página Processamento "Exatamente uma vez" e confirmação "No máximo uma vez" para mais informações.
- Duração da função com valor de tabela:evite durações muito curtas e excessivamente longas para funções com valor de tabela. Recomendamos uma duração moderada, como 20 minutos.
- Dimensionamento de lote:ajuste
max_batch_sizeà sua carga de trabalho. Use lotes menores para eventos de distribuição de dados alta e evite a disputa de bloqueio em linhas compartilhadas. Use lotes maiores para tarefas independentes com muitas consultas por segundo. Confirme ou estenda os períodos de concessão de mensagens para um lote em uma única transação para ter o melhor desempenho.
Monitoramento
É possível monitorar operações de fila usando tabelas de introspecção do Spanner. Embora essas tabelas não incluam colunas específicas da fila, é possível identificar a atividade da fila pesquisando os nomes definidos pelo usuário nas seguintes tabelas:
- Estatísticas de leitura
- Estatísticas de transação
- Estatísticas de bloqueio
- Estatísticas de tamanhos de tabela
- Estatísticas de operações de tabela
As seguintes métricas de fila do Spanner são encontradas no prefixo spanner.googleapis.com/queue/* no Cloud Monitoring:
buffered_ready_messages: (GAUGE, INT64, 1) a contagem de mensagens mantidas na memória e prontas para entrega a um receptor.message_send_count: (DELTA, INT64, 1) o número de mensagens enviadas no Spanner durante o intervalo de uma fila.message_ack_count: (DELTA, INT64, 1) O número de mensagens confirmadas no Spanner durante o intervalo de uma fila.oldest_unacked_message_age: (GAUGE, INT64, 1) a idade (em segundos) da mensagem mais antiga não confirmada em uma fila.lease_expiration_count: (DELTA, INT64, 1) o número de expirações de concessão no Spanner durante o intervalo de uma fila.
Todas as métricas anteriores são amostradas aproximadamente a cada 60 segundos. Após a amostragem, os dados podem não ficar visíveis por até 120 segundos. O registro de auditoria do Spanner abrange gravações, leituras e operações de esquema em filas.
Resolver problemas
As seções a seguir descrevem como identificar e resolver problemas comuns ao usar filas do Spanner.
O backlog de mensagens está crescendo
Diagnóstico
oldest_unacked_message_age e buffered_ready_messages estão elevados. Isso indica um desequilíbrio entre a taxa de envio de mensagens e a capacidade de processamento de mensagens do seu aplicativo.
Resolução
Para resolver esse problema, faça o seguinte:
- Verifique a taxa de confirmação:se a métrica
message_ack_countcaiu, verifique os workers do cliente para garantir que eles estejam funcionando corretamente e não tenham parado ou falhado. - Verifique a taxa de envio:se
message_send_counttiver aumentado, escalonar verticalmente dos workers de processamento de mensagens para lidar com o aumento da carga. - Confirme o esgotamento de recursos:verifique se as contagens de
buffered_ready_messageselease_expiration_countestão elevadas. Essa combinação indica uma escassez de receptores de função com valor de tabela (TVF) ativos ou um processamento lento do cliente.
As mensagens individuais estão travadas
Diagnóstico
A métrica oldest_unacked_message_age está alta, mas buffered_ready_messages está baixa ou estável. Isso indica que as mensagens individuais não estão sendo processadas ou confirmadas, e não um gargalo geral de capacidade.
Resolução
Para resolver esse problema, faça o seguinte:
Identifique mensagens bloqueadas:consulte a tabela de filas e ordene os resultados por tempo de entrega para encontrar as mensagens não confirmadas mais antigas:
GoogleSQL
SELECT * FROM UserTasks ORDER BY DeliverTime ASC LIMIT 10;PostgreSQL
SELECT * FROM usertasks ORDER BY deliver_time ASC LIMIT 10;Investigue falhas de processamento:verifique os registros do aplicativo para determinar por que os workers não estão reconhecendo as mensagens identificadas. As mensagens não confirmadas são reenviadas automaticamente após o vencimento do contrato.
Os contratos de mensagens expiram antes da conclusão do processamento
Diagnóstico
A métrica lease_expiration_count está elevada ou aumentando. Isso indica que o tempo de processamento da mensagem excede a duração do lease (padrão de 10 segundos) antes que os workers possam confirmar as mensagens.
Resolução
Para resolver esse problema, faça o seguinte:
- Renove os contratos de locação de forma proativa:se o processamento da mensagem levar mais de 10
segundos, chame a
TVF
RENEWLEASE_QUEUE_NAME()periodicamente. Renove o lease aproximadamente 7 a 8 segundos após o início do processamento para considerar a latência de rede com segurança. - Investigue o processamento lento:se o aplicativo já renova concessões
ativamente, mas
lease_expiration_countcontinua alto, verifique o código do back-end para identificar gargalos de processamento, chamadas de RPC lentas ou impasses.
Não é possível dimensionar a taxa de envio de mensagens
Diagnóstico
A métrica message_send_count atinge um platô de capacidade de processamento ou as solicitações de publicação encontram uma latência de gravação elevada quando você tenta aumentar a taxa de envio.
Resolução
Para resolver esse problema, considere as seguintes mudanças estruturais:
- Escalonar verticalmente os recursos de computação:adicione nós ou unidades de processamento à sua instância do Spanner para aumentar a capacidade geral do banco de dados.
- Verifique as divisões:avalie se é possível adicionar mais divisões para distribuir a carga de gravação entre vários servidores.
- Otimize o roteamento:verifique se o aplicativo de publicação grava diretamente na região líder da instância do Spanner para minimizar a latência de gravação.
As mensagens prontas não estão sendo processadas
Diagnóstico
A métrica buffered_ready_messages está alta e aumentando. Isso indica que
as mensagens estão em buffer e prontas para entrega na memória, mas os workers receptores
não estão as extraindo.
Resolução
Para resolver esse problema, faça o seguinte:
- Verifique as conexões TVF ativas:confira a contagem de TVF do receptor ativo para garantir que os workers do leitor estejam se conectando e extraindo mensagens ativamente. Se os workers forem desconectados ou não executarem consultas
RECEIVE_QUEUE_NAME()simultâneas suficientes, as mensagens vão permanecer sem sondagem no buffer.
A seguir
- Confira mais cenários e exemplos de filas do Spanner.
- 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.