En este documento, se proporcionan patrones de arquitectura y ejemplos de código para situaciones comunes de mensajería con colas de Spanner. Puedes usar estos patrones para activar el trabajo asíncrono después de que se confirmen las transacciones, programar tareas recurrentes o retrasadas, administrar cargas útiles de mensajes grandes con almacenamiento fuera de banda, coordinar flujos de trabajo de varios eventos y extender o crear puntos de control para trabajos en segundo plano de larga duración.
Procesamiento “exactamente una vez” y confirmación “a lo sumo una vez”
En la página Procesamiento de tipo “exactamente una vez” y confirmación de recepción de tipo “como máximo una vez”, se describen con más detalle las diversas consideraciones y soluciones para el procesamiento de tipo “exactamente una vez” y la confirmación de recepción de tipo “como máximo una vez”.
Realiza el trabajo después de que se confirma una transacción
Para realizar el trabajo después de que se confirma una transacción, envía un mensaje a la cola dentro de la misma transacción.
Por ejemplo, el registro de un usuario nuevo activa un correo electrónico de bienvenida:
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)
);
Después de que se confirma la transacción, el receptor de UserTasks transmite el mensaje, envía el correo electrónico y confirma la recepción del mensaje:
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';
Cómo controlar el trabajo de larga duración
Si tienes trabajo que podría tardar más que el arrendamiento predeterminado (más de 10 segundos), llama a SELECT * FROM RENEWLEASE_QUEUE_NAME() de forma periódica.
Por ejemplo, para generar un informe, haz lo siguiente:
- El receptor recibe un mensaje de
RECEIVE_ReportQueue(). - Iniciar la generación del informe
- Cada 5 segundos, llama a
SELECT * FROM RENEWLEASE_ReportQueue([leaseToken])en un subproceso o una rutina independientes. - Cuando se complete, confirma la recepción del mensaje y almacena el informe.
Como alternativa, si tienes trabajo de larga duración que requiere un procesamiento de, como máximo, una vez, o un tiempo de concesión prolongado, haz lo siguiente:
- Reconoce (
DELETEoACK) el mensaje actual de la fila al llegar. En la misma transacción, vuelve a poner en cola un mensaje nuevo con una marca de tiempo de entrega en el futuro, más allá del tiempo que lleva el procesamiento. - Continúa con el procesamiento y confirma el mensaje recién puesto en cola cuando termines.
Las ventajas de este enfoque son que no es necesario extender el arrendamiento de forma continua y que el mensaje no se vuelve a entregar hasta que llega la hora futura (lo que abarca las fallas). Si la confirmación inicial se realiza correctamente, se logra un procesamiento de, como máximo, una vez.
Crea puntos de control para el trabajo de larga duración
Las filas de Spanner pueden administrar tareas que duran de minutos a horas, no solo trabajos rápidos. Para estas tareas de larga duración, usa el siguiente enfoque:
- Almacena metadatos de forma externa: Usa almacenamiento fuera de banda para guardar los detalles y el estado de la tarea.
- Crea puntos de control con regularidad: Para recuperarse de fallas sin perder mucho progreso, la tarea debe guardar su estado de forma periódica.
- Usa el patrón de confirmación recomendado: La mejor manera de confirmar es reconocer de forma atómica (
ACK) el mensaje actual de la cola y enviar un mensaje nuevo programado para su entrega en el futuro. Este nuevo mensaje contiene o apunta al estado actualizado, lo que evita la reentrega inmediata a otro trabajador.
Este patrón reduce el trabajo duplicado, incluso si no es posible crear puntos de control completos, aunque la tarea se reinicia desde el principio después de una falla en ese caso.
Programar trabajo para un momento específico en el futuro
Para programar el trabajo para un momento específico en el futuro, establece la columna DeliverTime cuando insertes el mensaje.
Por ejemplo, un recordatorio de vencimiento de la prueba:
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');
Cómo controlar cargas útiles de mensajes grandes
Si la carga útil del mensaje es grande, usa el patrón de almacenamiento fuera de banda. Almacena la carga útil grande en una tabla separada y coloca una referencia a ella en el mensaje de la cola.
Por ejemplo, el procesamiento de imágenes:
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.
Espera varios eventos antes de continuar
Para esperar varios eventos antes de continuar (como una operación de unión), usa una tabla para hacer un seguimiento del estado y una cola para activar las verificaciones.
Por ejemplo, el cumplimiento de pedidos que requiere inventario y pago:
- Crea una tabla
OrdersconInventoryStatusyPaymentStatus. - Cuando se confirme el inventario, actualiza
Ordersy envía un mensaje aOrderCheckQueue. - Cuando se confirme el pago, actualiza
Ordersy envía un mensaje aOrderCheckQueue. - El receptor de
OrderCheckQueueverifica la tablaOrders. Si se confirman ambos estados, se procede con el envío y se confirma la recepción del mensaje. De lo contrario, es posible que se vuelva a poner en cola para una revisión posterior o que se ejecute otra lógica.
Cómo realizar una acción periódicamente
Para realizar una acción de forma periódica, usa el patrón de programación periódica. El receptor confirma la recepción del mensaje y envía uno nuevo programado para el siguiente intervalo.
Por ejemplo, la agregación de datos 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, usa las mutaciones Ack y Send de la biblioteca cliente. En estos ejemplos, se supone que tienes un objeto Message que encapsula la clave y la carga útil:
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();
});
}
¿Qué sigue?
- Aprende a usar las colas de Spanner, incluidas las prácticas recomendadas y la supervisión.
- 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.