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:
- Los tokens de arrendamiento que no se pueden analizar no devuelven una fila.
- Los tokens de arrendamiento que ya vencieron no devuelven una fila.
- Los tokens de alquiler no renovables sí devuelven una fila con un
SpannerNewLeaseTokende 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_sizea 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:
- Lee estadísticas
- Estadísticas de transacciones
- Estadísticas de bloqueo
- Estadísticas de tamaños de tablas
- Estadísticas de operaciones de 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_countdisminuyó, 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_countaumentó 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_messagesylease_expiration_countson 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_countsigue 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?
- Explora más situaciones y ejemplos de colas de Spanner.
- Obtén más información sobre el procesamiento de tipo “exactamente una vez” y la confirmación de recepción de tipo “como máximo una vez”.
- Configura el control de acceso con el control de acceso detallado para las colas.