Saltar al contenido
Open Security
Backend y ArquitecturaIntermedio· 45 min

NATS y JetStream: subjects, streams y consumers sin magia

Publicar no es lo mismo que persistir. Aprendé Core NATS, JetStream, wildcards, streams y consumers con un laboratorio local.

#backend#mensajeria#nats

Antes de empezar necesitás

  • Docker funcionando en tu máquina
  • NATS CLI instalado (nats)
  • jq instalado para leer JSON con comodidad
  • Saber abrir una terminal y ejecutar comandos

Al terminar vas a poder

  • Diferenciar Core NATS de JetStream
  • Publicar y consumir mensajes usando subjects
  • Usar wildcards * y > sin confundir publicación con suscripción
  • Crear un stream de JetStream y verificar qué mensajes persiste
  • Entender qué hace un consumer y por qué los ACK importan
  • Mapear variables de entorno típicas de una app que publica y consume eventos

Cuando dos servicios tienen que hablar, hay una pregunta incómoda: ¿tienen que estar conectados al mismo tiempo?

Si la respuesta es sí, Core NATS alcanza muchas veces. Si la respuesta es no, necesitás persistencia, replay y ACKs. Ahí entra JetStream.

Este lab no es para memorizar comandos. Es para que puedas mirar una configuración real de mensajería y responder tres cosas:

  1. Dónde se publica un evento.
  2. Qué stream lo guarda.
  3. Quién lo consume y desde qué subject.

1. NATS en una frase

NATS es un sistema de mensajería basado en subjects. Un servicio publica un mensaje en un subject, y otro servicio se suscribe a ese subject para recibirlo.

vt@labs:~
nats pub events.created '{"id":"evt-1"}'
nats sub events.created

El subject es el punto de encuentro. No necesitás saber la IP del consumidor, ni cuántas instancias hay, ni dónde están corriendo. Publicás en un nombre lógico.

Core NATS sirve muy bien para:

  • Request/reply entre servicios.
  • Fan-out a consumidores conectados.
  • Comunicación de baja latencia.
  • Señales o comandos que no necesitás reproducir después.

2. JetStream en una frase

JetStream es la capa de persistencia de NATS. Guarda mensajes que matchean ciertos subjects dentro de un stream, y después los entrega a través de consumers.

La diferencia práctica es esta:

Core NATS:     publico -> si alguien escucha, recibe
JetStream:     publico -> el stream guarda -> un consumer entrega y espera ACK

JetStream sirve cuando necesitás:

  • No perder eventos si el consumidor está caído.
  • Reprocesar mensajes desde el inicio, desde una hora o desde una secuencia.
  • Escalar workers con ACKs y redelivery.
  • Tener una cola de trabajo persistente.
  • Auditar qué eventos pasaron por un flujo.

3. Tecnologías parecidas

No todas resuelven exactamente el mismo problema, pero compiten en la misma zona:

  • Apache Kafka: fuerte para event streaming, alto throughput y ecosistema grande. Más pesado de operar.
  • RabbitMQ: broker clásico de colas, routing flexible y AMQP. Muy usado para workloads tradicionales.
  • Redis Streams: streams persistentes dentro de Redis. Simple si ya usás Redis, pero con otro modelo operativo.
  • Apache Pulsar: streaming distribuido con separación entre cómputo y storage. Potente, más complejo.
  • MQTT brokers: muy usados en IoT. El foco está en dispositivos, conexiones inestables y topics livianos.

4. Antes de empezar

Para este lab necesitás tres cosas en tu máquina:

  • Docker: para levantar un servidor NATS descartable.
  • NATS CLI (nats): para publicar, suscribirte y administrar JetStream.
  • jq: para leer los payloads JSON sin romperte los ojos.

Si no tenés la CLI de NATS, la podés instalar así:

vt@labs:~
# macOS con Homebrew
brew install nats-io/nats-tools/nats

# Linux (descarga directa, revisá la última versión en https://github.com/nats-io/natscli)
curl -sf https://get-nats.io | sh

# Verificá que esté disponible
nats --version

5. Levantá NATS con JetStream

Corré un servidor local descartable. No uses IPs ni tokens reales para este lab.

vt@labs:~
docker run --rm -p 4222:4222 -p 8222:8222 nats:latest -js -m 8222
  • -js activa JetStream.
  • -m 8222 expone el endpoint HTTP de monitoreo.
  • --rm hace que el contenedor se borre al frenarlo.

En otra terminal, guardá un contexto local para no repetir el server en cada comando:

vt@labs:~
nats context save local-lab --server=nats://127.0.0.1:4222 --select
nats context ls

6. Subjects y wildcards

NATS organiza los subjects por tokens separados con punto:

orders
orders.new
orders.new.created
orders.payment.created

Hay dos wildcards importantes:

  • * matchea un solo token.
  • > matchea uno o más tokens y solo puede ir al final.

Ejemplos:

orders.*     matchea orders.new
orders.*     no matchea orders.new.created

orders.>     matchea orders.new
orders.>     matchea orders.new.created

Probalo con dos suscriptores. En una terminal:

terminal 1
nats sub 'orders.*'

En otra terminal:

terminal 2
nats pub orders.new '{"type":"one-token"}'
nats pub orders.new.created '{"type":"two-tokens"}'

El subscriber orders.* debería ver solo orders.new.

Ahora escuchá todo lo que cuelga de orders:

terminal 1
nats sub 'orders.>'

Y publicá otra vez:

terminal 2
nats pub orders.new '{"type":"one-token"}'
nats pub orders.new.created '{"type":"two-tokens"}'

Ahora deberían llegar los dos.

7. El detalle que rompe configuraciones

Este punto parece menor, pero rompe integraciones reales:

orders.> matchea orders.new
orders.> no matchea orders

Un stream configurado solo con orders.> no captura mensajes publicados al subject exacto orders. Para aceptar ambos, configurá los dos subjects:

orders, orders.>

Creá un stream que capture ambos casos:

vt@labs:~
nats stream add ORDERS --subjects "orders,orders.>" --defaults
nats stream ls

Publicá tres mensajes:

vt@labs:~
nats pub orders '{"event":"root-subject"}'
nats pub orders.payment.created '{"event":"payment.created","orderId":"ord-1"}'
nats pub orders.shipment.created '{"event":"shipment.created","orderId":"ord-1"}'

Miralos desde el stream:

vt@labs:~
nats stream view ORDERS --translate "jq ."

Si querés acotar por tiempo:

vt@labs:~
nats stream view ORDERS --since 1h --translate "jq ."

8. El flujo completo, en un diagrama

Cuando publicás un mensaje, JetStream decide si lo guarda mirando los subjects del stream. Después un consumer lo entrega bajo sus propias reglas.

flowchart LR
  P[Publisher] -->|publica en| S(subject concreto)
  S -->|matchea| ST[Stream ORDERS]
  ST -->|entrega| C[Consumer orders-worker]
  C -->|ACK| ST
Publisher -> subject -> stream -> consumer. Cada flecha es un lugar donde puede fallar la configuración.

Ese diagrama es la cadena que tenés que poder recorrer cuando un mensaje “no llega”:

  1. ¿El publisher publicó en el subject correcto?
  2. ¿El stream captura ese subject?
  3. ¿El consumer filtra ese subject?
  4. ¿El consumer está activo y confirma los ACKs?

9. Streams no son subjects

Un subject es el nombre lógico donde publicás o escuchás mensajes.

Un stream es el almacenamiento de JetStream que captura mensajes publicados en uno o más subjects.

Un stream puede capturar varios subjects:

Stream: ORDERS
Subjects: orders, orders.>

Entonces, cuando publicás esto:

subject: orders.payment.created
payload: {"orderId":"ord-1"}

JetStream revisa si algún stream captura ese subject. Como ORDERS tiene orders.>, lo guarda.

10. Consumers: quién lee y confirma

El stream guarda mensajes. El consumer decide cómo se entregan.

Un consumer mantiene estado:

  • Qué mensajes ya entregó.
  • Qué mensajes fueron confirmados con ACK.
  • Qué mensajes debe redeliverar si no hubo ACK.
  • Desde qué punto empieza a leer.

Creá un consumer durable para leer solo eventos bajo orders.>:

vt@labs:~
nats consumer add ORDERS orders-worker --filter "orders.>" --ack explicit --pull --defaults

Pedile el próximo mensaje:

vt@labs:~
nats consumer next ORDERS orders-worker

Pull vs push

Los consumers pueden ser de dos tipos:

  • Push: JetStream empuja los mensajes al cliente cuando llegan. Es útil para latencia baja, pero el cliente debe estar conectado y listo para recibir.
  • Pull: el cliente pide mensajes cuando quiere. Es más fácil de escalar con múltiples workers y evita que un consumidor lento se sature.

En este lab usamos --pull porque es más fácil de experimentar desde la terminal con nats consumer next.

11. Request/reply sigue siendo Core NATS

No todo tiene que pasar por JetStream. Si necesitás una respuesta inmediata, Core NATS tiene request/reply.

En una terminal levantá un responder:

terminal 1
nats reply rpc.create_request '{"requestId":"req-123","projectId":"proj-123"}'

En otra terminal mandá una request:

terminal 2
nats request rpc.create_request '{"action":"create"}'

Ese subject no tiene por qué estar en un stream. Es una conversación síncrona: si se completa, listo; si no hay responder o hay timeout, falla.

12. Cómo leer una configuración real

En una app real, el problema casi nunca es NATS “en abstracto”. El problema es descubrir qué valores terminan formando el subject final.

Un ejemplo saneado:

NATS_ENABLED=true
NATS_URL=nats://127.0.0.1:4222

JETSTREAM_ENABLED=true
JETSTREAM_STREAM_NAME=ORDERS
JETSTREAM_STREAM_PUBLISH_NAME=ORDERS
JETSTREAM_STREAM_CONSUMER_NAME=ORDERS

JETSTREAM_SUBJECTS_PUBLISH=shop.orders.events,shop.orders.events.>
JETSTREAM_SUBJECT_PREFIX=shop.orders.events

ORDERS_TASK_SUBJECT=shop.orders.events.>
DLQ_SUBJECT=dlq.>

Hay tres preguntas que tenés que hacerle a esa configuración:

  1. ¿En qué stream publico?
  2. ¿Con qué prefix se arma el subject final?
  3. ¿Qué subject está escuchando el consumer?

Si una app publica así:

await js.publish(
    subject=f"{settings.NATS_SUBJECT_PREFIX}.{subject}",
    payload=nats_event.model_dump_json().encode(),
    stream=settings.NATS_STREAM_NAME,
)

Y NATS_SUBJECT_PREFIX=shop.orders.events, entonces publicar payment.created termina en:

shop.orders.events.payment.created

El consumidor tiene que escuchar algo que matchee eso, por ejemplo:

shop.orders.events.>

Si escucha orders.> o shop.orders.>, no va a recibir ese evento.

13. Limpieza

Cuando termines, borrá el consumer y el stream para no dejar datos dando vueltas:

vt@labs:~
nats consumer rm ORDERS orders-worker --force
nats stream rm ORDERS --force

Después frená el contenedor con Ctrl+C en la terminal donde lo levantaste. Como usaste --rm, el contenedor desaparece solo.

14. Checklist de debugging

Cuando no llega un mensaje, no arranques cambiando código. Seguí el recorrido.

Primero, verificá el contexto:

vt@labs:~
nats context ls
nats rtt

Después, verificá el stream:

vt@labs:~
nats stream ls
nats stream info ORDERS

Miralo con JSON legible:

vt@labs:~
nats stream view ORDERS --since 1h --translate "jq ."

Y recién ahí revisá la app:

  • Subject final que publica.
  • Stream configurado para publicar.
  • Subjects que captura el stream.
  • Filter subject del consumer.
  • ACKs pendientes o redeliveries.
  • DLQ o subject de error si existe.

Cerrando

NATS no es una cola mágica. Es subject-based messaging.

JetStream no es otro broker separado. Es la persistencia de NATS: streams para guardar, consumers para entregar y ACKs para saber qué pasó.

Si te llevás una sola idea, que sea esta:

Publisher -> subject concreto -> stream captura por subject -> consumer entrega con reglas

Cuando algo no llega, buscá dónde se cortó esa cadena.

Lo que practicás en este lab

Llevátelo a tu repo si querés, pero no es obligatorio: es tu aprendizaje.

  • Salida de nats stream ls mostrando el stream creado
  • Salida de nats stream view con al menos dos eventos JSON
  • Writeup corto explicando la diferencia entre subject, stream y consumer
  • Ejemplo de configuración saneada sin IPs privadas, tokens ni datos reales

Reto

Creá un stream TASKS que capture tasks y tasks.>, publicá tres eventos JSON, consumilos con un consumer durable y explicá por qué tasks.> no captura el subject exacto tasks. Para probar la validación interactiva, ingresá la flag: flag{jetstream_persists}

¿Hiciste el lab?

Si querés, guardá lo que hiciste (comandos, notas, un repo) para volver después. Y si encontrás un error o querés mejorar este lab,contribuí al repo. El progreso se guarda solo en tu navegador.