Empresas
Empleos
  • Sobre nosotros
  • Soluciones
    • Publicación de vacantes
      Publica tu vacante y recibe candidatos calificados en 48h.
    • Evaluación de candidatos
      500+ pruebas técnicas y psicológicas, más anti-fraude.
    • Headhunting
      Búsqueda ejecutiva a la medida de principio a fin.
    • Nómina + EOR
      Dispersión de nómina y EOR en más de 15 países de LATAM.
  • Precios
  • Empleos

0

284
Vistas
Mongo change-Stream con Spring resumeAt vs startAfter y tolerancia a fallas en caso de pérdida de conexión

No puedo encontrar una respuesta en stackOverflow, ni en ninguna documentación, tengo el siguiente código de flujo de cambios (escuche un DB, no una colección específica)

La versión de Mongo es 4.2

 @Configuration public class DatabaseChangeStreamListener { //Constructor, fields etc... @PostConstruct public void initialize() { MessageListenerContainer container = new DefaultMessageListenerContainer(mongoTemplate, new SimpleAsyncTaskExecutor(), this::onException); ChangeStreamRequest.ChangeStreamRequestOptions options = new ChangeStreamRequest.ChangeStreamRequestOptions(mongoTemplate.getDb().getName(), null, buildChangeStreamOptions()); container.register(new ChangeStreamRequest<>(this::onDatabaseChangedEvent, options), Document.class); container.start(); } private ChangeStreamOptions buildChangeStreamOptions() { return ChangeStreamOptions.builder() .returnFullDocumentOnUpdate() .filter(newAggregation(match(where(OPERATION_TYPE).in(INSERT.getValue(), UPDATE.getValue(), REPLACE.getValue(), DELETE.getValue())))) .resumeAt(Instant.now().minusSeconds(1)) .build(); } //more code }

Quiero que la transmisión comience a escuchar solo desde el momento de inicio del sistema, sin tomar nada antes en el registro de operaciones, ¿ .resumeAt(Instant.now().minusSeconds(1)) ? ¿Necesito usar el método starAfter Si es así, ¿cómo puedo encontrar el último resumeToken en la base de datos? ¿O está listo para usar y no necesito agregar ningún currículum/líneas de inicio?

segunda pregunta, nunca detengo el contenedor (siempre debe vivir mientras se ejecuta la aplicación), en caso de desconexión de mongoDB y reconexión, ¿el oyente en la configuración actual continuará consumiendo mensajes? (Estoy teniendo dificultades para simular la desconexión de la base de datos)

Si no reanuda el manejo de eventos, ¿qué debo cambiar en la configuración para que el flujo de cambios continúe y tome todos los eventos del último resumeToken recibido antes de la desconexión? He leído este excelente artículo sobre el flujo de cambio medio en producción , pero usa el cursor directamente, y quiero usar el Spring DefaultMessageListenerContainer , ya que es mucho más elegante.

over 4 years ago · Santiago Trujillo
1 Respuestas
Responde la pregunta

0

Así que responderé mis propias preguntas (algunas más tontas, otras menos :)...):

  1. cuando no se proporciona la marca de tiempo resumeAt, la secuencia de cambios comenzará a partir de la hora actual y no dibujará ningún evento anterior.
  2. La diferencia entre el evento resumeAfter y la marca de tiempo se puede encontrar aquí: respuesta de stackOverflow , pero tenga en cuenta que para la marca de tiempo incluye el evento, por lo que si desea comenzar desde el próximo evento (en Java), haga lo siguiente:
 private BsonTimestamp getNextEventTimestamp(BsonTimestamp timestamp) { return new BsonTimestamp(timestamp.getValue() + 1); }
  1. En caso de desconexión de Internet, el flujo de cambios no se reanudará, por lo que recomiendo seguir el siguiente enfoque en caso de error:
 private void onException() { ScheduledExecutorService executorService = newSingleThreadScheduledExecutor(); executorService.scheduleAtFixedRate(() -> recreateChangeStream(executorService), 0, 1, TimeUnit.SECONDS); } private void recreateChangeStream(ScheduledExecutorService executorService) { try { mongoTemplate.getDb().runCommand(new BasicDBObject("ping", "1")); container.stop(); startNewContainer(); executorService.shutdown(); } catch (Exception ignored) { } }

Primero, estoy creando una tarea programada ejecutable que siempre se ejecuta (pero solo 1 a la vez newSingleThreadScheduledExecutor() ), estoy tratando de hacer ping a la base de datos, después de un ping exitoso, detengo el contenedor anterior y comienzo uno nuevo, también puede pase la última marca de tiempo que tomó para que pueda obtener todos los eventos que podría haberse perdido

recuperación de la marca de tiempo del evento:

 BsonTimestamp resumeAtTimestamp = changeStreamDocument.getClusterTime();

entonces estoy cerrando la tarea.

también asegúrese de que resumeAtTimestamp exista en oplog ...

over 4 years ago · Santiago Trujillo Denunciar
Responde la pregunta
Encuentra empleos remotos

¡Descubre la nueva forma de encontrar empleo!

Top de empleos
Top categorías de empleo
Empresas
Publicar vacante Precios Comercial
Legal
Términos y condiciones Política de privacidad
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomiéndame algunas ofertas
Necesito ayuda