Redis Stream — la fiabilidad y escalabilidad de sus sistemas de mensajería

Redis Stream — la fiabilidad y escalabilidad de sus sistemas de mensajería

Redis Stream: un nuevo tipo de dato abstracto introducido en Redis con la versión 5.0
Conceptualmente, Redis Stream es una lista en la que puedes agregar registros. Cada registro tiene un identificador único. Por defecto, el identificador se genera automáticamente e incluye una marca de tiempo. Por lo tanto, puedes solicitar rangos de registros por tiempo o recibir nuevos datos a medida que llegan al flujo, como el comando Unix «tail -f» lee un archivo de registro y se detiene esperando nuevos datos. Ten en cuenta que múltiples clientes pueden escuchar el flujo simultáneamente, de la misma manera que muchos procesos «tail -f» pueden leer un archivo al mismo tiempo sin conflictos.

Para entender todas las ventajas del nuevo tipo de dato, repasemos brevemente las estructuras de Redis que ya existen y que parcialmente replican la funcionalidad de Redis Stream.

Redis PUB/SUB

Redis Pub/Sub es un sistema de mensajería simple, ya integrado en tu almacenamiento clave-valor. Sin embargo, a cambio de su simplicidad, se deben considerar las siguientes desventajas:

  • Si el publicador falla por alguna razón, pierde todos sus suscriptores.
  • El publicador necesita conocer la dirección exacta de todos sus suscriptores.
  • El publicador puede sobrecargar a sus suscriptores si se publican datos más rápido de lo que pueden procesarlos.
  • El mensaje se elimina del búfer del publicador inmediatamente después de su publicación, independientemente de cuántos suscriptores lo recibieron y con qué rapidez pudieron procesarlo.
  • Todos los suscriptores recibirán el mensaje al mismo tiempo. Los suscriptores deben coordinar entre ellos el orden de procesamiento del mismo mensaje.
  • No hay un mecanismo integrado para confirmar la exitosa recepción de un mensaje por parte del suscriptor. Si el suscriptor recibe el mensaje y falla durante su procesamiento, el publicador no lo sabrá.

Redis List

Redis List es una estructura de datos que soporta comandos de lectura con bloqueo. Puedes añadir y leer mensajes desde el principio o el final de la lista. Con esta estructura, puedes construir una buena pila o cola para tu sistema distribuido y esto en la mayoría de los casos será suficiente. Las principales diferencias con Redis Pub/Sub son:

  • El mensaje se entrega a un solo cliente. El primer cliente que bloquea la lectura recibirá los datos primero.
  • Clint debe iniciar la operación de lectura de cada mensaje por sí mismo. List no sabe nada sobre los clientes.
  • Los mensajes se almacenan hasta que alguien los cuenta o los elimina explícitamente. Si has configurado el servidor Redis para que realice un volcado de datos en disco, la confiabilidad del sistema aumenta drásticamente.

Introducción a Stream

Agregar un registro al stream

Comando XADD añade un nuevo registro al stream. Un registro no es solo una cadena, consiste en uno o más pares clave-valor. Por lo tanto, cada registro ya está estructurado y recuerda la estructura de un archivo CSV.

> XADD mystream * sensor-id 1234 temperature 19.8
1518951480106-0

En el ejemplo anterior, estamos añadiendo al stream llamado (clave) 'mystream' dos campos: 'sensor-id' y 'temperature' con los valores '1234' y '19.8' respectivamente. Como segundo argumento, el comando toma un identificador que se asignará al registro; este identificador identifica de manera única cada registro en el stream. Sin embargo, en este caso hemos pasado *, porque queremos que Redis genere un nuevo identificador para nosotros. Cada nuevo identificador se incrementará. Por lo tanto, cada nuevo registro tendrá un identificador mayor en relación con los registros anteriores.

Formato del identificador

El identificador del registro, devuelto por el comando XADD, consiste en dos partes:

{millisecondsTime}-{sequenceNumber}

millisecondsTime — tiempo Unix en milisegundos (tiempo servidores Redis). Sin embargo, si el tiempo actual resulta ser igual o menor que el tiempo del registro anterior, se utiliza la marca de tiempo del registro anterior. Por lo tanto, si el tiempo del servidor regresa al pasado, el nuevo identificador aún tendrá que mantener la propiedad de incremento.

sequenceNumber se utiliza para los registros creados en el mismo milisegundo. sequenceNumber se incrementará en 1 con respecto al registro anterior. Dado que sequenceNumber tiene un tamaño de 64 bits, en la práctica no deberías alcanzar un límite en la cantidad de registros que pueden generarse en un milisegundo.

El formato de estos identificadores puede parecer extraño a primera vista. Un lector desconfiado puede preguntarse por qué el tiempo es parte del identificador. La razón es que los flujos de Redis admiten consultas de rango por identificadores. Dado que el identificador está relacionado con el tiempo de creación del registro, esto permite realizar consultas sobre rangos de tiempo. Veremos un ejemplo específico cuando pasemos a estudiar el comando. XRANGE.

Si por alguna razón el usuario necesita especificar su propio identificador, que, por ejemplo, esté relacionado con algún sistema externo, podemos pasarlo al comando XADD en lugar del símbolo * como se muestra a continuación:

> XADD somestream 0-1 field value
0-1
> XADD somestream 0-2 foo bar
0-2

Tenga en cuenta que en este caso debe hacerse cargo del incremento del identificador. En nuestro ejemplo, el identificador mínimo es "0-1", por lo que el comando no aceptará otro identificador que sea igual o menor que "0-1".

> XADD somestream 0-1 foo bar
(error) ERR El ID especificado en XADD es igual o menor que el elemento superior del flujo objetivo.

Número de registros en el flujo

Se puede obtener el número de registros en el flujo simplemente usando el comando XLEN. Para nuestro ejemplo, este comando devolverá el siguiente valor:

> XLEN somestream
(integer) 2

Consultas de rango — XRANGE y XREVRANGE

Para solicitar datos por rango, necesitamos especificar dos identificadores: el inicio y el final del rango. El rango devuelto incluirá todos los elementos, incluyendo los límites. También hay dos identificadores especiales «-» y «+», que representan respectivamente el identificador más pequeño (el primer registro) y el más grande (el último registro) en el flujo. El ejemplo a continuación mostrará todos los registros del flujo.

> XRANGE mystream - +
1) 1) 1518951480106-0
   2) 1) "sensor-id"
      2) "1234"
      3) "temperature"
      4) "19.8"
2) 1) 1518951482479-0
   2) 1) "sensor-id"
      2) "9999"
      3) "temperature"
      4) "18.2"

Cada registro devuelto representa un array de dos elementos: el identificador y una lista de pares clave-valor. Ya hemos mencionado que los identificadores de los registros están relacionados con el tiempo. Por lo tanto, podemos solicitar un rango de un intervalo de tiempo específico. Sin embargo, podemos especificar en la consulta no un identificador completo, sino solo el tiempo Unix, omitiendo la parte relacionada con sequenceNumberLa parte omitida del identificador se asignará automáticamente a cero al comienzo del rango y al valor máximo posible al final del rango. A continuación se muestra un ejemplo de cómo solicitar un rango de dos milisegundos.

> XRANGE mystream 1518951480106 1518951480107
1) 1) 1518951480106-0
   2) 1) "sensor-id"
      2) "1234"
      3) "temperatura"
      4) "19.8"

Solo tenemos un registro en este rango, sin embargo, en conjuntos de datos reales, el resultado devuelto puede ser enorme. Por esta razón, XRANGE se admite la opción COUNT. Al especificar una cantidad, podemos simplemente obtener los primeros N registros. Si necesitamos obtener los siguientes N registros (paginación), podemos usar el último identificador obtenido, incrementarlo sequenceNumber en uno y solicitar nuevamente. Veamos esto en el siguiente ejemplo. Comenzamos a agregar 10 elementos usando XADD (supongamos que el flujo mystream ya se ha llenado con 10 elementos). Para comenzar la iteración, obteniendo 2 elementos por comando, comenzamos con el rango completo, pero con COUNT igual a 2.

> XRANGE mystream - + COUNT 2
1) 1) 1519073278252-0
   2) 1) "foo"
      2) "value_1"
2) 1) 1519073279157-0
   2) 1) "foo"
      2) "value_2"

Para continuar la iteración con los siguientes dos elementos, necesitamos seleccionar el último identificador obtenido, es decir, 1519073279157-0, y agregar 1 a sequenceNumber.
el identificador resultante, en este caso 1519073279157-1, que ahora puede usarse como un nuevo argumento de inicio del rango para la siguiente llamada. XRANGE:

> XRANGE mystream 1519073279157-1 + COUNT 2
1) 1) 1519073280281-0
   2) 1) "foo"
      2) "value_3"
2) 1) 1519073281432-0
   2) 1) "foo"
      2) "value_4"

Y así sucesivamente. Dado que la complejidad XRANGE es O(log (N)) para la búsqueda y luego O(M) para devolver M elementos, cada paso de la iteración es rápido. Así, con XRANGE se puede iterar eficientemente sobre los flujos.

Comando XREVRANGE es el equivalente XRANGE, pero devuelve elementos en orden inverso:

> XREVRANGE mystream + - COUNT 1
1) 1) 1519073287312-0
   2) 1) "foo"
      2) "value_10"

Tenga en cuenta que el comando XREVRANGE toma los argumentos del rango de inicio y parada en orden inverso.

Leer nuevos registros usando XREAD

A menudo surge la necesidad de suscribirse a un flujo y recibir solo mensajes nuevos. Este concepto puede parecerse a Redis Pub/Sub o a la lista bloqueante de Redis, pero existen diferencias fundamentales en cómo usar Redis Stream:

  1. Cada nuevo mensaje se entrega por defecto a cada suscriptor. Este comportamiento es diferente de la lista bloqueante de Redis, donde un nuevo mensaje solo será leído por un único suscriptor.
  2. Mientras que en Redis Pub/Sub todos los mensajes se olvidan y nunca se almacenan, en Stream todos los mensajes se mantienen indefinidamente (a menos que el cliente solicite explícitamente la eliminación).
  3. Redis Stream permite restringir el acceso a los mensajes dentro de un mismo stream. Un suscriptor específico puede ver solo su propio historial de mensajes.

Puedes suscribirte al stream y recibir nuevos mensajes usando el comando XREAD. Esto es un poco más complicado que XRANGE, así que primero comenzaremos con ejemplos más simples.

> XREAD COUNT 2 STREAMS mystream 0
1) 1) "mystream"
   2) 1) 1) 1519073278252-0
         2) 1) "foo"
            2) "value_1"
      2) 1) 1519073279157-0
         2) 1) "foo"
            2) "value_2"

En el ejemplo anterior se indica la forma no bloqueante XREAD. Ten en cuenta que la opción COUNT no es obligatoria. De hecho, la única opción obligatoria del comando es la opción STREAMS, que especifica la lista de streams junto con el identificador máximo correspondiente. Escribimos “STREAMS mystream 0” — queremos recibir todos los registros del stream mystream con el identificador mayor que “0-0”. Como se puede ver en el ejemplo, el comando devuelve el nombre del stream, porque podemos suscribirnos a varios streams al mismo tiempo. Podríamos escribir, por ejemplo, “STREAMS mystream otherstream 0 0”. Ten en cuenta que después de la opción STREAMS primero necesitamos proporcionar los nombres de todos los streams necesarios y solo luego la lista de identificadores.

En esta forma simple, el comando no hace nada especial en comparación con XRANGE. Sin embargo, lo interesante es que podemos convertir fácilmente XREAD en un comando bloqueante, especificando el argumento BLOCK:

> XREAD BLOCK 0 STREAMS mystream $

En el ejemplo anterior, se ha indicado una nueva opción BLOCK con un tiempo de espera de 0 milisegundos (esto significa espera infinita). Además, en lugar de pasar un identificador normal para el stream mystream, se ha pasado un identificador especial $. Este identificador especial significa que XREAD debe usar el identificador máximo en el stream mystream. Así que solo recibiremos nuevos mensajes a partir del momento en que comenzamos a escuchar. En cierto sentido, esto es parecido al comando Unix «tail -f».

Tenga en cuenta que al usar la opción BLOCK no es necesario utilizar un identificador especial $. Podemos usar cualquier identificador existente en el flujo. Si el comando puede atender nuestra solicitud de inmediato, sin bloquear, lo hará; de lo contrario, se bloqueará.

Bloqueante XREAD también puede escuchar varios flujos a la vez, solo es necesario especificar sus nombres. En este caso, el comando devolverá el registro del primer flujo en el que se recibieron los datos. El primer suscriptor que sea bloqueado para dicho flujo recibirá los datos primero.

Grupos de Consumidores

En ciertas tareas, queremos restringir el acceso de los suscriptores a los mensajes dentro de un solo flujo. Un ejemplo de cuándo puede ser útil esto es en una cola de mensajes con trabajadores que recibirán diferentes mensajes del flujo, permitiendo escalar el procesamiento de mensajes.

Si imaginamos que tenemos tres suscriptores C1, C2, C3 y un flujo que contiene los mensajes 1, 2, 3, 4, 5, 6, 7, el servicio de mensajes ocurrirá como se muestra en el diagrama a continuación:

1 -> C1
2 -> C2
3 -> C3
4 -> C1
5 -> C2
6 -> C3
7 -> C1

Para lograr este efecto, Redis Stream utiliza un concepto llamado Grupo de Consumidores. Este concepto es similar a un pseudo-suscriptor que recibe datos del flujo, pero en realidad es atendido por varios suscriptores dentro del grupo, proporcionando ciertas garantías:

  1. Cada mensaje se entrega a suscriptores diferentes dentro del grupo.
  2. Dentro del grupo, los suscriptores se identifican por un nombre que es una cadena que distingue entre mayúsculas y minúsculas. Si algún suscriptor se cae temporalmente del grupo, puede reingresar al grupo con su propio nombre único.
  3. Cada Grupo de Consumidores sigue el concepto de "primer mensaje no leído". Cuando un suscriptor solicita nuevos mensajes, solo puede recibir aquellos que nunca han sido entregados a ningún suscriptor dentro del grupo.
  4. Hay un comando que confirma explícitamente el procesamiento exitoso del mensaje por el suscriptor. Mientras no se ejecute este comando, el mensaje solicitado permanecerá en estado "pendiente".
  5. Dentro del Grupo de Consumidores, cada suscriptor puede solicitar el historial de los mensajes que se le han entregado, pero que aún no han sido procesados (en estado "pendiente").

En cierto sentido, el estado del grupo puede representarse así:

+----------------------------------------+
| consumer_group_name: mygroup          
| consumer_group_stream: somekey        
| last_delivered_id: 1292309234234-92    
|                                                           
| consumers:                                          
|    "consumer-1" con mensajes pendientes  
|       1292309234234-4                          
|       1292309234232-8                          
|    "consumer-42" con mensajes pendientes 
|       ... (y así sucesivamente)                             
+----------------------------------------+

Ahora es momento de familiarizarse con los comandos principales para el Grupo de Consumidores, a saber:

  • XGROUP se utiliza para crear, destruir y gestionar grupos
  • XREADGROUP se utiliza para leer flujos a través del grupo
  • XACK — este comando permite al suscriptor marcar un mensaje como procesado con éxito

Creando un Grupo de Consumidores

Supongamos que el flujo mystream ya existe. Entonces, el comando para crear el grupo será:

> XGROUP CREATE mystream mygroup $
OK

Al crear el grupo debemos pasar el identificador desde el cual el grupo comenzará a recibir mensajes. Si solo queremos recibir todos los nuevos mensajes, podemos usar un identificador especial $ (como en nuestro ejemplo anterior). Si en lugar de un identificador especial usamos 0, el grupo tendrá acceso a todos los mensajes del flujo.

Ahora que el grupo está creado, podemos empezar a leer los mensajes de inmediato usando el comando XREADGROUP. Este comando es muy parecido a XREAD y soporta la opción opcional BLOCK. Sin embargo, hay una opción obligatoria GROUP, que siempre debe especificarse con dos argumentos: el nombre del grupo y el nombre del suscriptor. La opción COUNT también es soportada.

Antes de leer el flujo, coloquemos algunos mensajes allí:

> XADD mystream * message apple
1526569495631-0
> XADD mystream * message orange
1526569498055-0
> XADD mystream * message strawberry
1526569506935-0
> XADD mystream * message apricot
1526569535168-0
> XADD mystream * message banana
1526569544280-0

Y ahora intentemos leer este flujo a través del grupo:

> XREADGROUP GROUP mygroup Alice COUNT 1 STREAMS mystream >
1) 1) "mystream"
   2) 1) 1) 1526569495631-0
         2) 1) "message"
            2) "apple"

El comando anterior dice literalmente lo siguiente:

«Yo, Alice, suscriptor, miembro del grupo mygroup, quiero leer del flujo mystream un mensaje que nunca ha sido entregado a nadie antes».

Cada vez que un suscriptor realiza una operación con un grupo, debe indicar su nombre, identificándose de manera única dentro del grupo. En el comando anterior, hay otro detalle muy importante: un identificador especial ">". Este identificador especial filtra los mensajes, dejando solo aquellos que aún no se han entregado en ninguna ocasión.

Además, en casos especiales, puede especificar un identificador real, como 0 o cualquier otro identificador válido. En este caso, el comando XREADGROUP te devolverá el historial de mensajes con estado «pending», que han sido entregados al suscriptor especificado (Alice), pero que aún no han sido confirmados mediante el comando XACK.

Podemos verificar este comportamiento, especificando de inmediato el identificador 0, sin la opción COUNT. Solo veremos un único mensaje pendiente, es decir, el mensaje con la manzana:

> XREADGROUP GROUP mygroup Alice STREAMS mystream 0
1) 1) "mystream"
   2) 1) 1) 1526569495631-0
         2) 1) "message"
            2) "apple"

Sin embargo, si confirmamos el mensaje como procesado con éxito, ya no aparecerá:

> XACK mystream mygroup 1526569495631-0
(integer) 1
> XREADGROUP GROUP mygroup Alice STREAMS mystream 0
1) 1) "mystream"
   2) (lista o conjunto vacío)

Ahora es el turno de Bob de leer algo:

> XREADGROUP GROUP mygroup Bob COUNT 2 STREAMS mystream >
1) 1) "mystream"
   2) 1) 1) 1526569498055-0
         2) 1) "message"
            2) "orange"
      2) 1) 1526569506935-0
         2) 1) "message"
            2) "strawberry"

Bob, un miembro del grupo mygroup, pidió no más de dos mensajes. El comando solo informa sobre los mensajes no entregados debido al identificador especial ">". Como puedes ver, el mensaje «apple» no aparece, ya que ya fue entregado a Alice, por lo que Bob recibe «orange» y «strawberry».

Así, Alice, Bob y cualquier otro suscriptor del grupo pueden leer diferentes mensajes del mismo flujo. También pueden leer su historial de mensajes no procesados o marcar mensajes como procesados.

Hay varias cosas que tener en cuenta:

  • Una vez que un suscriptor lee un mensaje con el comando XREADGROUP, este mensaje pasa a un estado de «pending» y se asigna a este suscriptor en particular. Otros suscriptores del grupo no podrán leer este mensaje.
  • Los suscriptores se crean automáticamente en la primera mención, no es necesario crearlos explícitamente.
  • Con XREADGROUP puedes leer mensajes de diferentes flujos simultáneamente, sin embargo, para que esto funcione, necesitas crear previamente grupos con el mismo nombre para cada flujo usando XGROUP

Recuperación de fallos

Un suscriptor puede recuperarse de un fallo y volver a leer su lista de mensajes con el estado 'pendiente'. Sin embargo, en el mundo real los suscriptores pueden fallar definitivamente. ¿Qué ocurre con los mensajes de suscripción pendientes si no pudo recuperarse del fallo?
El Consumer Group ofrece una función que se utiliza precisamente para tales casos — cuando es necesario cambiar la propiedad de los mensajes.

Lo primero que debes hacer es llamar al comando XPENDING, que muestra todos los mensajes del grupo con el estado 'pendiente'. En su forma más simple, el comando se llama solo con dos argumentos: el nombre del flujo y el nombre del grupo:

> XPENDING mystream mygroup
1) (integer) 2
2) 1526569498055-0
3) 1526569506935-0
4) 1) 1) "Bob"
      2) "2"

El comando devolvió la cantidad de mensajes no procesados para todo el grupo y para cada suscriptor. Solo tenemos a Bob con dos mensajes no procesados, ya que el único mensaje solicitado por Alice fue confirmado con XACK.

Podemos solicitar información adicional utilizando más argumentos:

XPENDING {key} {groupname} [{start-id} {end-id} {count} [{consumer-name}]]

{start-id} {end-id} — rango de identificadores (se pueden usar '-' y '+')
{count} — número de intentos de entrega
{consumer-name} — nombre del grupo

> XPENDING mystream mygroup - + 10
1) 1) 1526569498055-0
   2) "Bob"
   3) (integer) 74170458
   4) (integer) 1
2) 1) 1526569506935-0
   2) "Bob"
   3) (integer) 74170458
   4) (integer) 1

Ahora tenemos detalles para cada mensaje: identificador, nombre del suscriptor, tiempo de espera en milisegundos y, finalmente, el número de intentos de entrega. Tenemos dos mensajes de Bob, y llevan 74170458 milisegundos pendientes, aproximadamente 20 horas.

Tenga en cuenta que nada nos impide verificar cuál fue el contenido del mensaje, simplemente usando XRANGE.

> XRANGE mystream 1526569498055-0 1526569498055-0
1) 1) 1526569498055-0
   2) 1) "mensaje"
      2) "naranja"

Solo necesitamos repetir el mismo identificador dos veces en los argumentos. Ahora que tenemos una idea, Alice puede decidir que después de 20 horas de espera, Bob probablemente no se recuperará, y es hora de solicitar estos mensajes y reanudar su procesamiento en lugar de Bob. Para esto usamos el comando XCLAIM:

XCLAIM {key} {group} {consumer} {min-idle-time} {ID-1} {ID-2} ... {ID-N}

Con este comando podemos obtener un mensaje "ajeno" que aún no ha sido procesado, cambiando el propietario a {consumer}. Sin embargo, también podemos proporcionar un tiempo de inactividad mínimo de {min-idle-time}. Esto ayuda a evitar situaciones en las que dos clientes intentan cambiar simultáneamente el propietario de los mismos mensajes.

Cliente 1: XCLAIM mystream mygroup Alice 3600000 1526569498055-0
Cliente 2: XCLAIM mystream mygroup Lora 3600000 1526569498055-0

El primer cliente restablecerá el tiempo de inactividad y aumentará el contador de entregas. Por lo tanto, el segundo cliente no podrá solicitarlo.

> XCLAIM mystream mygroup Alice 3600000 1526569498055-0
1) 1) 1526569498055-0
   2) 1) "mensaje"
      2) "naranja"

El mensaje ha sido reclamado con éxito por Alice, quien ahora puede procesar el mensaje y confirmarlo.

Del ejemplo anterior, podemos ver que la ejecución exitosa de la solicitud devuelve el contenido del mensaje en sí. Sin embargo, esto no es obligatorio. La opción JUSTID se puede utilizar para devolver solo los identificadores del mensaje. Esto es útil si no te interesan los detalles del mensaje y deseas mejorar el rendimiento del sistema.

Contador de entregas

El contador que observas en la salida XPENDING es la cantidad de entregas de cada mensaje. Este contador se incrementa de dos maneras: cuando el mensaje se reclama con éxito a través de XCLAIM o cuando se utiliza la llamada XREADGROUP.

Es normal que algunos mensajes sean entregados varias veces. Lo principal es que, al final, todos los mensajes sean procesados. A veces, al procesar un mensaje, pueden surgir problemas debido a la corrupción del propio mensaje o el procesamiento del mensaje puede causar un error en el código del manejador. En tal caso, puede resultar que nadie pueda procesar este mensaje. Dado que tenemos un contador de intentos de entrega, podemos usar este contador para detectar tales situaciones. Por lo tanto, una vez que el contador de entregas alcance un número grande que hayas establecido, probablemente sea más sensato mover tal mensaje a otro flujo y enviar una notificación al administrador del sistema.

Estado de los flujos

Comando XINFO se utiliza para solicitar diversa información sobre el flujo y sus grupos. Por ejemplo, la forma básica del comando es la siguiente:

> XINFO STREAM mystream
 1) longitud
 2) (entero) 13
 3) claves-radix-tree
 4) (entero) 1
 5) nodos-radix-tree
 6) (entero) 2
 7) grupos
 8) (entero) 2
 9) primera-entrada
10) 1) 1524494395530-0
    2) 1) "a"
       2) "1"
       3) "b"
       4) "2"
11) última-entrada
12) 1) 1526569544280-0
    2) 1) "mensaje"
       2) "banana"

El comando anterior muestra información general sobre el flujo indicado. Ahora, un ejemplo un poco más complicado:

> XINFO GROUPS mystream
1) 1) nombre
   2) "mygroup"
   3) consumidores
   4) (entero) 2
   5) pendiente
   6) (entero) 2
2) 1) nombre
   2) "some-other-group"
   3) consumidores
   4) (entero) 1
   5) pendiente
   6) (entero) 0

El comando anterior muestra información general sobre todos los grupos del flujo indicado.

> XINFO CONSUMERS mystream mygroup
1) 1) nombre
   2) "Alice"
   3) pendiente
   4) (entero) 1
   5) inactivo
   6) (entero) 9104628
2) 1) nombre
   2) "Bob"
   3) pendiente
   4) (entero) 1
   5) inactivo
   6) (entero) 83841983

El comando anterior muestra información sobre todos los suscriptores del flujo y del grupo indicado.
Si olvidas la sintaxis del comando, simplemente pide ayuda al propio comando:

> XINFO HELP
1) XINFO {subcomando} arg arg ... arg. Los subcomandos son:
2) CONSUMERS {clave} {nombre del grupo}  -- Muestra los grupos de consumidores del grupo {nombre del grupo}.
3) GROUPS {clave}                 -- Muestra los grupos de consumidores del flujo.
4) STREAM {clave}                 -- Muestra información sobre el flujo.
5) HELP                         -- Imprime esta ayuda.

Límite del tamaño del flujo

Muchas aplicaciones no quieren acumular datos en un flujo para siempre. A menudo es útil tener una cantidad máxima permitida de mensajes en el flujo. En otros casos, es útil mover todos los mensajes del flujo a otro almacenamiento permanente al alcanzar un tamaño de flujo especificado. El tamaño del flujo se puede limitar usando el parámetro MAXLEN en el comando. XADD:

> XADD mystream MAXLEN 2 * value 1
1526654998691-0
> XADD mystream MAXLEN 2 * value 2
1526654999635-0
> XADD mystream MAXLEN 2 * value 3
1526655000369-0
> XLEN mystream
(entero) 2
> XRANGE mystream - +
1) 1) 1526654999635-0
   2) 1) "value"
      2) "2"
2) 1) 1526655000369-0
   2) 1) "value"
      2) "3"

Al usar MAXLEN, los registros antiguos se eliminan automáticamente al alcanzar la longitud especificada, por lo que el flujo tiene un tamaño constante. Sin embargo, el recorte en este caso no se realiza de la manera más eficiente en la memoria de Redis. La situación se puede mejorar de la siguiente manera:

XADD mystream MAXLEN ~ 1000 * ... campos de entrada aquí ...

El argumento ~ en el ejemplo anterior significa que no es necesario limitar la longitud del flujo a un valor específico. En nuestro ejemplo, esto puede ser cualquier número mayor o igual a 1000 (por ejemplo, 1000, 1010 o 1030). Simplemente hemos indicado que queremos que nuestro flujo mantenga al menos 1000 registros. Esto hace que el manejo de la memoria dentro de Redis sea mucho más eficiente.

También existe un comando separado XTRIM, que realiza lo mismo:

> XTRIM mystream MAXLEN 10

> XTRIM mystream MAXLEN ~ 10

Almacenamiento permanente y replicación

Redis Stream se replica de forma asíncrona en nodos esclavos y se almacena en archivos tipo AOF (instantánea de todos los datos) y RDB (registro de todas las operaciones de escritura). También se admite la replicación del estado de los Grupos de Consumidores. Por lo tanto, si un mensaje está en estado «pendiente» en el nodo maestro, en los nodos esclavos este mensaje tendrá el mismo estado.

Eliminación de elementos individuales del flujo

Para eliminar mensajes existe un comando especial XDEL. El comando recibe el nombre del flujo, seguido de los identificadores de los mensajes que necesitan ser eliminados:

> XRANGE mystream - + COUNT 2
1) 1) 1526654999635-0
   2) 1) "valor"
      2) "2"
2) 1) 1526655000369-0
   2) 1) "valor"
      2) "3"
> XDEL mystream 1526654999635-0
(integer) 1
> XRANGE mystream - + COUNT 2
1) 1) 1526655000369-0
   2) 1) "valor"
      2) "3"

Al usar este comando, hay que tener en cuenta que la memoria no se liberará de inmediato.

Flujos de longitud cero

La diferencia entre los flujos y otras estructuras de datos de Redis es que, cuando otras estructuras de datos ya no tienen elementos dentro, como efecto secundario, la propia estructura de datos se eliminará de la memoria. Por ejemplo, un conjunto ordenado se eliminará por completo cuando el llamado a ZREM elimine el último elemento. En cambio, se permite que los flujos permanezcan en la memoria incluso si no tienen ningún elemento dentro.

Conclusión

Redis Stream es ideal para crear mensajeros, colas de mensajes, registros unificados y sistemas de chat que almacenan historial.

Como una vez dijo Niklaus Wirth, los programas son algoritmos más estructuras de datos, y Redis ya te proporciona ambos.

Fuente: habr.com

Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS 🔥 Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS | ProHoster