Usa colas de Spanner

En este documento, se describe cómo usar las colas de Spanner. Explica cómo crear una cola, enviar y recibir mensajes, extender los arrendamientos de mensajes y confirmar la recepción de mensajes. También incluye prácticas recomendadas, información sobre las colas de supervisión y orientación para solucionar problemas.

Crea una cola

Para crear una cola, usa la instrucción 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;

La única columna que no es necesario crear de forma explícita en la instrucción CREATE QUEUE se llama DeliverTime en GoogleSQL y deliver_time en PostgreSQL. Spanner los crea automáticamente.

Las filas admiten políticas de tiempo de actividad (TTL), que pueden ayudar a administrar el backlog de mensajes antiguos que no se confirmaron.

Envía un mensaje

Para enviar un mensaje a una cola, usa la declaración 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');

También puedes usar las mutaciones de la biblioteca cliente Send para insertar un mensaje:

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

Recibir mensajes

Usa la función con valor de tabla (TVF) RECEIVE_QUEUE_NAME() con ExecuteStreamingSQL para recibir mensajes. Esta es una llamada de larga duración. Debes ejecutar una de estas llamadas por trabajador y por fila de forma repetitiva. Ejecuta la consulta con una lectura sólida porque Spanner rechaza las lecturas inactivas.

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

Tu código de cliente debe iterar los resultados con una consulta de transmisión. Cada fila devuelta es un mensaje.

Cómo recibir mensajes en un lote

Para aumentar la capacidad de procesamiento y procesar varios mensajes juntos, puedes recibir mensajes en lotes especificando el 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');

Como alternativa, usa la biblioteca cliente. En este ejemplo de Go, se muestra cómo transmitir mensajes desde una cola, verificar el vencimiento del alquiler y confirmar mensajes de forma así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)
}

Extiende el arrendamiento del mensaje

Si el procesamiento de un mensaje tarda más que la concesión inicial (10 segundos), usa la siguiente sintaxis:

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"

Los tokens de arrendamiento se devuelven según la siguiente lógica:

  1. Los tokens de arrendamiento que no se pueden analizar no devuelven una fila.
  2. Los tokens de arrendamiento que ya vencieron no devuelven una fila.
  3. Los tokens de alquiler no renovables devuelven una fila con un SpannerNewLeaseToken de NULL. Esto puede ocurrir si el mensaje ya se confirmó, pero el token de arrendamiento no venció.

Como alternativa, usa la biblioteca cliente. En este ejemplo de Go, se muestra cómo extender el arrendamiento de un mensaje:

// 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
}

Cómo confirmar la recepción de un mensaje

Usa el DELETE DML para confirmar la recepción de un mensaje. Debe ser transaccional con cualquier otra escritura relacionada con el procesamiento de mensajes.

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, usa la mutación Ack de la biblioteca cliente para confirmar la recepción de un mensaje:

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ácticas recomendadas

A continuación, se indican las prácticas recomendadas para usar las colas de Spanner:

  • Cargas útiles pequeñas: Mantén las cargas útiles de los mensajes de la cola pequeñas, por debajo de los 4 KB. Usa el patrón de almacenamiento fuera de banda para datos más grandes.
  • Administración de arrendamientos: Extiende los arrendamientos de las tareas que podrían superar el arrendamiento de mensajes predeterminado de 10 segundos. Si no se extienden los arrendamientos de mensajes, se pueden producir reenvíos y un posible procesamiento doble.
  • Manejo de errores: Spanner pone en cola los mensajes que no se pueden procesar y los reintenta con una espera exponencial durante la primera hora. Los mensajes más antiguos se reintentarán una vez por hora. Considera mover los mensajes que fallan de forma permanente a una cola separada.
  • Supervisión: Supervisa las profundidades de las colas y la antigüedad del mensaje sin confirmar más antiguo para asegurarte de que tus receptores no limiten la entrada de tu canalización.
  • Idempotencia: Diseña tus procesadores de mensajes para que toleren operaciones repetidas, ya que la entrega al menos una vez significa que es posible que se envíen duplicados ocasionalmente. Consulta la página Procesamiento de tipo exactamente una vez y confirmación de recepción de tipo a lo sumo una vez para obtener más información.
  • Duración de la función con valores de tabla: Evita duraciones extremadamente cortas y excesivamente largas para las funciones con valores de tabla. Se recomienda una duración moderada, como 20 minutos.
  • Tamaño del lote: Ajusta max_batch_size a tu carga de trabajo. Usa lotes más pequeños para los eventos con una gran cantidad de fan-out para evitar la contención de bloqueos en las filas compartidas. Usa lotes más grandes para tareas independientes con una gran cantidad de consultas por segundo. Confirma o extiende los arrendamientos de mensajes para un lote en una sola transacción para obtener el mejor rendimiento.

Supervisar

Puedes supervisar las operaciones de la cola con las tablas de introspección de Spanner. Si bien estas tablas no incluyen columnas específicas de la cola, puedes identificar la actividad de la cola buscando los nombres de las colas definidos por el usuario en las siguientes tablas:

Las siguientes métricas de la cola de Spanner se encuentran bajo el prefijo spanner.googleapis.com/queue/* en Cloud Monitoring:

  • buffered_ready_messages: (GAUGE, INT64, 1) Es el recuento de mensajes que se almacenan en la memoria y están listos para entregarse a un receptor.
  • message_send_count: (DELTA, INT64, 1) Es la cantidad de mensajes enviados en Spanner durante el intervalo para una cola.
  • message_ack_count: (DELTA, INT64, 1) Es la cantidad de mensajes confirmados en Spanner durante el intervalo para una cola.
  • oldest_unacked_message_age: (GAUGE, INT64, 1) Es la antigüedad (en segundos) del mensaje no confirmado más antiguo en una cola.
  • lease_expiration_count: (DELTA, INT64, 1) Es la cantidad de vencimientos de arrendamiento en Spanner durante el intervalo para una cola.

Todas las métricas anteriores se muestrean aproximadamente cada 60 segundos. Después del muestreo, es posible que los datos no se puedan ver durante un máximo de 120 segundos. El registro de auditoría de Spanner abarca las operaciones de escritura, lectura y esquema en las colas.

Solucionar problemas

En las siguientes secciones, se describe cómo identificar y resolver problemas comunes cuando se usan las colas de Spanner.

La cantidad de mensajes pendientes está aumentando

Diagnóstico

Tanto oldest_unacked_message_age como buffered_ready_messages están elevados. Esto indica un desequilibrio entre la tasa de envío de mensajes y la capacidad de procesamiento de mensajes de tu aplicación.

Solución

Para solucionar este problema, haz lo siguiente:

  • Verifica la tasa de confirmación: Si la métrica message_ack_count disminuyó, verifica tus trabajadores del cliente para asegurarte de que se ejecuten correctamente y no se hayan detenido ni fallado.
  • Verifica la tasa de envío: Si message_send_count aumentó repentinamente, escala verticalmente la cantidad de trabajadores de procesamiento de mensajes para controlar el aumento de la carga.
  • Confirma el agotamiento de recursos: Verifica si los recuentos de buffered_ready_messages y lease_expiration_count son elevados. Esta combinación indica una escasez de receptores activos de funciones con valor de tabla (TVF) o un procesamiento lento del cliente.

Los mensajes individuales se atascan

Diagnóstico

La métrica oldest_unacked_message_age es alta, pero la métrica buffered_ready_messages es baja o estable. Esto indica que los mensajes individuales no se procesan ni se confirman, en lugar de que haya un cuello de botella general en la capacidad.

Solución

Para solucionar este problema, haz lo siguiente:

  • Identifica los mensajes atascados: Consulta la tabla de la cola y ordena los resultados por hora de entrega para encontrar los mensajes sin confirmar más antiguos:

    GoogleSQL

    SELECT *
    FROM UserTasks
    ORDER BY DeliverTime ASC
    LIMIT 10;
    

    PostgreSQL

    SELECT *
    FROM usertasks
    ORDER BY deliver_time ASC
    LIMIT 10;
    
  • Investiga las fallas en el procesamiento: Revisa los registros de tu aplicación para determinar por qué los trabajadores no confirman la recepción de los mensajes identificados. Los mensajes no confirmados se vuelven a entregar automáticamente después de que vence su concesión.

Los arrendamientos de mensajes vencen antes de que se complete el procesamiento

Diagnóstico

La métrica lease_expiration_count es alta o está en aumento. Esto indica que el tiempo de procesamiento de mensajes supera la duración de la asignación (10 segundos de forma predeterminada) antes de que los trabajadores puedan confirmar los mensajes.

Solución

Para solucionar este problema, haz lo siguiente:

  • Renueva los arrendamientos de forma proactiva: Si el procesamiento de mensajes tarda más de 10 segundos, llama al RENEWLEASE_QUEUE_NAME() TVF periódicamente. Renueva el arrendamiento aproximadamente entre 7 y 8 segundos después de que comience el procesamiento para tener en cuenta de forma segura la latencia de la red.
  • Investiga el procesamiento lento: Si tu aplicación ya renueva las concesiones de forma activa, pero lease_expiration_count sigue siendo alto, revisa tu código de backend para detectar cuellos de botella en el procesamiento, llamadas a RPC lentas o interbloqueos.

No se puede ajustar la tasa de envío de mensajes

Diagnóstico

La métrica message_send_count alcanza una meseta de capacidad de procesamiento o las solicitudes de publicación experimentan una latencia de escritura elevada cuando intentas aumentar la tasa de envío.

Solución

Para resolver este problema, considera los siguientes cambios estructurales:

  • Escala verticalmente los recursos de procesamiento: Agrega nodos o unidades de procesamiento a tu instancia de Spanner para aumentar la capacidad general de la base de datos.
  • Verifica las divisiones: Evalúa si puedes agregar más divisiones para distribuir la carga de escritura en varios servidores.
  • Optimiza el enrutamiento: Asegúrate de que tu aplicación de publicación escriba directamente en la región principal de tu instancia de Spanner para minimizar la latencia de escritura.

No se procesan los mensajes listos

Diagnóstico

La métrica buffered_ready_messages es alta y está en aumento. Esto indica que los mensajes se almacenan en búfer y están listos para la entrega en la memoria, pero los trabajadores receptores no los recuperan.

Solución

Para solucionar este problema, haz lo siguiente:

  • Verifica las conexiones de TVF activas: Verifica la cantidad de TVF de receptor activas para asegurarte de que los trabajadores lectores se conecten y extraigan mensajes de forma activa. Si los trabajadores se desconectaron o no ejecutan suficientes consultas RECEIVE_QUEUE_NAME() simultáneas, los mensajes permanecerán sin sondear en el búfer.

¿Qué sigue?