Nuestra aplicación frontend envía las acciones del usuario a una función lambda detrás de una puerta de enlace API, que luego almacena estas acciones en dynamodb.
Luego, usamos flujos de dynamodb para activar una función lambda separada que analizará estas acciones en dynamodb y decidirá si las acciones del usuario deben generar el envío de notificaciones (llamamos a estos eventos de notificación).
Por ejemplo, si un usuario coloca un comentario en nuestra aplicación, almacenaremos una acción "CREATED_COMMENT" en dynamodb, que luego activará una nueva lambda a través de una transmisión de dynamodb. La nueva lambda puede crear un "evento de notificación por correo electrónico", que podemos enviar a un proveedor de correo electrónico como customer.io
Sin embargo, nuestros usuarios nos informaron que reciben correos electrónicos con demasiada frecuencia y, por lo tanto, nos gustaría comenzar a enviar resúmenes de correos electrónicos agregando varias acciones a lo largo del tiempo en un solo correo electrónico en lugar de enviar un correo electrónico para cada acción.
Nuestra idea era usar algo como AWS EventBridge, Kinesis, Step Functions o incluso transmisiones de DynamoDB para reenviar las acciones de la transmisión de dynamodb, pero luego configurar los eventos de la nueva transmisión para que se agrupen por dirección de correo electrónico y para que estos eventos se repitan, por ejemplo, 10 minutos. Si el usuario luego realiza una nueva acción, la transmisión de ese usuario continuará recopilando acciones durante otros 10 minutos, hasta que no haya nuevas acciones de ese usuario durante 10 minutos. Una vez que eso suceda, la transmisión "liberará" todas las acciones reunidas e invocará una función lambda. Nuestra función lambda generará el evento de notificación por correo electrónico y lo enviará a, por ejemplo, customer.io.
Sin embargo, no hemos podido encontrar dicha agrupación y configuración de vaciado sin rebote en ninguno de los servicios de transmisión de AWS mencionados anteriormente. Para algo tan común como digerir (o resumir), ¿no debería haber un enfoque sin servidor para hacer esto sin tener que escribir nuestro propio servicio de cola?
La respuesta me parece usar una herramienta como SQS. SQS le permitirá acumular mensajes en una cola y cada x minutos podrá leer la cola mediante una función Lambda para hacerlo en un evento programado. No necesita tener un Lambda activado por SQS, y aún puede leer la cola "manualmente" desde dentro de Lambda.
Gareth McCumskey está en el camino correcto.
Use una cola normal de sqs estrictamente para eliminar rebotes.
Establezca una ventana de lote, es decir, 5 segundos. Use un tamaño de lote realmente grande cuando lea de la cola.
En el código, use un hashMap para agrupar su mensaje con el mismo ID de mensaje. Ahora use sus ID de mensajes deduplicados para hacer su trabajo.
Escribí una publicación de blog sobre algo como esto. La versión corta es que utiliza una función Lambda programada para identificar los registros que deben procesarse.
El problema con el uso de la demora en SQS es que solo puede recibir 10 mensajes a la vez, por lo que para recibir todos los mensajes, debe llamar a SQS repetidamente para borrar la cola. En ese momento, puede agregar los mensajes. Esto no escala muy bien, ya que todos los mensajes deben leerse para que funcione. Al usar DynamoDB, en realidad puede tener solo un registro que represente la colección de registros y consultar el registro único, lo que luego puede generar un mensaje en una cola para ese grupo específico de mensajes. Considere los siguientes datos:
user | comment | time user 1 | comment 1 | 11:43am user 1 | comment 2 | 11:50am user 2 | comment 1 | 11:51amPuede agregar otro registro que sea una señal de la necesidad de enviar un mensaje para cada usuario (en este ejemplo, 15 minutos después del primer mensaje).
user | scheduled user 1 | 11:58 user 2 | 12:06Cuando inserta el segundo conjunto de registros, está insertando la hora en que desea enviar el lote. Solo realiza la inserción si aún no hay un registro, por lo que no termina aumentando constantemente el tiempo. Su proceso programado lee ese registro para saber a qué usuarios necesita enviar mensajes y recopila todos los datos de ese usuario. El proceso de envío de mensajes a cada usuario se puede hacer en paralelo (podría enviar un mensaje al SQS para cada usuario o usar un estado de Mapa en una función de paso, por ejemplo).