MongoDB Academy · Lección

Reanudación de flujos de cambios tras una interrupción

Persistirá el token de reanudación y reiniciará un flujo de cambios desde el último evento procesado para garantizar una entrega al menos una vez.

Lección 4 de 413 pasos

Reanudación de flujos de cambios tras una interrupción es una lección gratuita de MongoDB Academy en CoddyKit. Esta es la lección 4 de 4. Puedes leer la lección completa abajo gratuitamente — luego la practicas en el navegador con un editor de código integrado y un tutor de IA 24/7. Forma parte de la ruta de aprendizaje de MongoDB Academy, y tu progreso se sincroniza en la web y la app de CoddyKit. El curso de MongoDB Academy incluye 4 lecciones en total.

El problema: ¿qué ocurre después de un fallo?

Un proceso consumidor de un flujo de cambios puede interrumpirse por fallos de red, despliegues, errores o reinicios intencionados. Sin un mecanismo para reanudarlo desde el último evento procesado, su aplicación reiniciaría el flujo desde el momento actual y perdería todos los eventos ocurridos durante la interrupción. Los tokens de reanudación de MongoDB resuelven este problema, ya que permiten volver a abrir un flujo desde una posición exacta del historial del oplog.

¿Qué es un token de reanudación?

Cada documento de evento de cambio tiene un campo _id que actúa como su token de reanudación: un valor binario opaco que identifica de forma única la posición del evento en el oplog. Tiene un aspecto como { _data: '...' }, donde los datos son un identificador interno codificado en hexadecimal. No necesita comprender su contenido; solo debe guardarlo y devolvérselo a MongoDB para reanudar el flujo desde esa posición exacta.

// A resume token is the _id field of a change event
// Example structure (your values will differ):
// {
//   _id: {
//     _data: '8264CE3B7200000002463C6F5A100...'
//   }
// }

// Store it as-is — do not parse or modify the _data value

Persistir el token de reanudación

Después de procesar cada evento, guarde el token de reanudación en un almacén persistente (base de datos, archivo o Redis) antes de confirmar que el evento se ha procesado. Este patrón de «checkpoint» garantiza que pueda continuar desde el último evento procesado correctamente después de un reinicio. Guárdelo después de procesar el evento, no antes, para evitar omitir eventos si el procesador falla durante el procesamiento.

const CHECKPOINT_COLLECTION = 'changeStreamCheckpoints';
const STREAM_ID = 'orderProcessorV1';

for await (const change of changeStream) {
  // Process the event
  await processOrderChange(change);

  // Save token AFTER successful processing
  await db.collection(CHECKPOINT_COLLECTION).updateOne(
    { streamId: STREAM_ID },
    { $set: { resumeToken: change._id, processedAt: new Date() } },
    { upsert: true }
  );
}

Reanudar con resumeAfter

Para reanudar desde un token guardado, páselo a watch() mediante la opción resumeAfter. MongoDB reproducirá los eventos posteriores al evento identificado por el token; el último evento procesado no se vuelve a reproducir. Si el token hace referencia a un evento que aún está en el oplog, MongoDB comenzará a entregar los eventos a partir del siguiente. Esto proporciona una entrega exactamente una vez cuando se combina con un procesamiento de eventos idempotente.

async function startOrResumeStream() {
  // Load the last saved resume token
  const checkpoint = await db.collection('changeStreamCheckpoints')
    .findOne({ streamId: 'orderProcessorV1' });

  const watchOptions = checkpoint
    ? { resumeAfter: checkpoint.resumeToken }
    : {}; // Start from current position if no token

  const changeStream = db.collection('orders').watch([], watchOptions);

  for await (const change of changeStream) {
    await processOrderChange(change);
    await saveResumeToken(change._id);
  }
}

startAfter frente a resumeAfter

MongoDB ofrece dos opciones similares: resumeAfter y startAfter. resumeAfter reanuda el flujo desde un token de operación normal, pero genera un error si recibe el token de un evento 'invalidate'. startAfter funciona de la misma manera, pero también puede reanudar el flujo después de un evento de invalidación, lo que resulta útil para volver a abrir un flujo sobre una colección después de que se haya eliminado y recreado. En la mayoría de los casos, resumeAfter es la opción adecuada.

// resumeAfter — standard use case, does not work after invalidate tokens
db.collection('orders').watch([], { resumeAfter: savedToken });

// startAfter — can also start after an invalidate event
// (use when the collection may have been dropped and recreated)
db.collection('orders').watch([], { startAfter: savedToken });

startAtOperationTime para reanudar según el tiempo

Si no tiene un token de reanudación, pero conoce la marca de tiempo desde la que desea comenzar, puede utilizar startAtOperationTime con un objeto Timestamp de MongoDB. Esto abre el flujo en un momento específico del clúster en lugar de hacerlo en un evento concreto. Resulta útil en escenarios de reproducción, por ejemplo, para volver a procesar todos los eventos desde el inicio de un despliegue, pero requiere que el oplog aún conserve el historial correspondiente a ese momento.

const { Timestamp } = require('mongodb');

// Replay events since a specific time
const startTime = new Timestamp({ t: Math.floor(Date.now() / 1000) - 3600, i: 1 }); // 1 hour ago

const changeStream = db.collection('orders').watch([], {
  startAtOperationTime: startTime
});

// Events from 1 hour ago will be delivered

Retención del oplog y el límite temporal

La reanudación de un flujo de cambios solo funciona si el evento correspondiente al token aún se encuentra en el oplog. El oplog de MongoDB es una colección limitada de tamaño finito; las entradas antiguas se sobrescriben a medida que llegan otras nuevas. En Atlas, puede configurar la retención del oplog, normalmente entre 24 y 72 horas. Si su aplicación estuvo inactiva durante más tiempo que el periodo de retención del oplog, el token apuntará a una posición caducada y MongoDB devolverá un error ChangeStreamHistoryLost (código 286). Su aplicación debe gestionar esta situación iniciando un flujo nuevo.

async function startOrResumeWithFallback() {
  const checkpoint = await loadResumeToken();
  try {
    const options = checkpoint ? { resumeAfter: checkpoint } : {};
    const stream = db.collection('orders').watch([], options);
    for await (const change of stream) {
      await processChange(change);
      await saveResumeToken(change._id);
    }
  } catch (error) {
    if (error.code === 286) { // ChangeStreamHistoryLost
      console.warn('Resume token expired — starting from now and running catch-up scan');
      await deleteResumeToken();
      await catchUpScan(); // full collection scan to catch missed changes
      await startOrResumeWithFallback();
    } else {
      throw error;
    }
  }
}

Procesamiento idempotente de eventos

Dado que los flujos de cambios proporcionan una entrega al menos una vez (el mismo evento podría entregarse más de una vez después de una reanudación), sus controladores de eventos deben ser idempotentes: procesar el mismo evento dos veces debe producir el mismo resultado que procesarlo una sola vez. Algunas técnicas son comprobar si el documento ya refleja el cambio antes de aplicarlo, utilizar upsert en lugar de insert o almacenar los ID de los eventos procesados en un conjunto de «eventos procesados» y omitir los duplicados.

async function processOrderChange(change) {
  if (change.operationType === 'insert') {
    const { fullDocument } = change;
    // Idempotent upsert — safe to replay:
    await db.collection('orderSummaries').updateOne(
      { _id: fullDocument._id },       // match by _id
      { $set: { ...summarize(fullDocument) } }, // idempotent set
      { upsert: true }                 // create if not exists
    );
  }
}

Reanudación automática a nivel del driver

El driver de MongoDB para Node.js (y otros drivers oficiales) incluye reanudación automática para errores de red transitorios. Cuando se pierde temporalmente la conexión con el servidor de MongoDB, el driver vuelve a abrir de forma transparente el flujo de cambios desde el token del último evento recibido, sin que su código tenga que gestionarlo. Esta reanudación automática cubre muchos escenarios habituales de interrupción; solo necesita implementar lógica de reanudación manual para reinicios a nivel de la aplicación (fallos del proceso o despliegues).

// The driver automatically resumes after network blips
// No code needed in your loop for this case:
for await (const change of changeStream) {
  // If the network drops and reconnects, the driver resumes automatically
  // and continues delivering events from where it left off.
  await processChange(change);
  await saveCheckpoint(change._id); // Still save tokens for process restarts
}

Controlador de flujos de cambios preparado para producción

Un consumidor de flujos de cambios preparado para producción debe combinar: persistencia del token de reanudación (para sobrevivir a los reinicios del proceso), gestión de la caducidad del oplog (con fallback a un escaneo completo), procesamiento idempotente (para reproducir eventos de forma segura) y apagado ordenado (para cerrar los flujos al recibir SIGTERM). Estos cuatro elementos juntos permiten procesar eventos de forma fiable incluso ante fallos de red, despliegues y periodos prolongados de inactividad.

class ReliableChangeStreamConsumer {
  constructor(collection, handler) {
    this.collection = collection;
    this.handler = handler;
    this.running = false;
  }

  async start() {
    this.running = true;
    while (this.running) {
      const token = await loadToken();
      const stream = this.collection.watch([], token ? { resumeAfter: token } : {});
      try {
        for await (const change of stream) {
          await this.handler(change);
          await saveToken(change._id);
        }
      } catch (e) {
        if (e.code === 286) { await deleteToken(); continue; } // expired, restart fresh
        if (!this.running) break; // shutting down
        throw e;
      } finally {
        await stream.close();
      }
    }
  }

  stop() { this.running = false; }
}

Apagado ordenado y vaciado del token

Al realizar un apagado planificado (implementación, reinicio), siempre vacíe el token de reanudación más reciente antes de cerrar el flujo de cambios. Si su aplicación procesa eventos en lotes (almacenándolos en búfer para mejorar el rendimiento), asegúrese de guardar el token después de cada lote, no después de cada evento individual. Realice un seguimiento del token del último evento procesado correctamente, no del último recibido. Esta distinción es importante: si recibe un evento, guarda su token y luego la aplicación falla antes de procesarlo, perderá ese evento al reanudar el flujo.

// Correct: save token AFTER processing, not before
for await (const change of changeStream) {
  await processChange(change);        // process first
  await saveToken(change._id);        // token saved only after success
  // If crash happens here, the event was already processed — safe
}

// Wrong: saving token before processing
for await (const change of changeStream) {
  await saveToken(change._id);        // token saved
  await processChange(change);        // if crash here, event is skipped on resume!
}

Comprobación rápida

Compruebe su comprensión de los conceptos de MongoDB y las bases de datos NoSQL de esta lección.

Resumen de la lección

En esta lección aprendió que el token de reanudación se almacena en change._id y debe persistirse después de procesar cada evento, que debe pasar el token a watch() como resumeAfter para reiniciar el flujo desde el último evento procesado y que debe gestionar ChangeStreamHistoryLost (código 286) cuando el token haya caducado en el oplog, reiniciando desde la hora actual y ejecutando un análisis de recuperación. A continuación, exploraremos cómo crear índices de Atlas Search para la búsqueda de texto completo.

Gratis para empezar

Aprende JavaScript con un tutor de IA — gratis

Escribe y ejecuta código real en tu navegador, obtén ayuda instantánea de un tutor de IA disponible 24/7 y continúa donde lo dejaste en la web o en la aplicación.

Cursos
30
Lecciones
120

Preguntas frecuentes

¿La lección «Reanudación de flujos de cambios tras una interrupción» es gratis?

Sí — el texto completo de «Reanudación de flujos de cambios tras una interrupción» es gratis para leer aquí en la web. Para practicarla de forma interactiva (editor de código integrado y tutor de IA 24/7) y desbloquear el resto del curso de MongoDB Academy, actualiza a CoddyKit PRO. El curso de MongoDB Academy incluye 4 lecciones en total.

¿Qué aprenderé en «Reanudación de flujos de cambios tras una interrupción»?

Persistirá el token de reanudación y reiniciará un flujo de cambios desde el último evento procesado para garantizar una entrega al menos una vez. Practicas MongoDB Academy con código real que ejecutas directamente en el navegador, y un tutor de IA 24/7 responde tus preguntas mientras trabajas en la lección.

¿Necesito experiencia previa para empezar MongoDB Academy?

No se requiere experiencia previa. MongoDB Academy en CoddyKit está estructurado para principiantes hasta estudiantes avanzados, así que puedes empezar aquí o desde el inicio y avanzar a tu ritmo. Esta es la lección 4 de 4.

¿Cuánto tiempo toma la lección «Reanudación de flujos de cambios tras una interrupción»?

La mayoría de las lecciones de CoddyKit toman alrededor de 5–10 minutos. Cada una es compacta e interactiva, así que avanzas constantemente y retomas exactamente por donde dejaste en la web y la app.

¿Puedo escribir y ejecutar código en esta lección de MongoDB Academy?

Sí. Cada lección de MongoDB Academy incluye un editor de código integrado, así que escribes y ejecutas código real directamente en tu navegador y obtienes retroalimentación instantánea de IA — sin configuración local necesaria.

Todas las lecciones de este curso

  1. Apertura de un flujo de cambios en una colección
  2. Estructura del documento de un evento de cambio
  3. Filtrado de eventos con un pipeline de agregación
  4. Reanudación de flujos de cambios tras una interrupción
← Volver a MongoDB Academy