Usar filas do Spanner

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:

  1. Tokens de concessão não analisáveis não retornam uma linha.
  2. Tokens de concessão já expirados não retornam uma linha.
  3. Tokens de concessão não renováveis não retornam uma linha com um SpannerNewLeaseToken de 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:

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_count caiu, 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_count tiver 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_messages e lease_expiration_count estã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_count continua 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