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.
Así que responderé mis propias preguntas (algunas más tontas, otras menos :)...):
private BsonTimestamp getNextEventTimestamp(BsonTimestamp timestamp) { return new BsonTimestamp(timestamp.getValue() + 1); } 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 ...