diff --git a/spring-5-data-reactive/src/main/java/com/baeldung/tailablecursor/service/InfoLogsCounter.java b/spring-5-data-reactive/src/main/java/com/baeldung/tailablecursor/service/InfoLogsCounter.java index b30eba0b25..29301bffec 100644 --- a/spring-5-data-reactive/src/main/java/com/baeldung/tailablecursor/service/InfoLogsCounter.java +++ b/spring-5-data-reactive/src/main/java/com/baeldung/tailablecursor/service/InfoLogsCounter.java @@ -18,8 +18,8 @@ public class InfoLogsCounter implements LogsCounter { public InfoLogsCounter(LogsRepository repository) { Flux stream = repository.findByLevel(LogLevel.INFO); - this.subscription = stream.subscribe(l -> { - log.info("INFO log received: " + l); + this.subscription = stream.subscribe(logEntity -> { + log.info("INFO log received: " + logEntity); counter.incrementAndGet(); }); } diff --git a/spring-5-data-reactive/src/main/java/com/baeldung/tailablecursor/service/WarnLogsCounter.java b/spring-5-data-reactive/src/main/java/com/baeldung/tailablecursor/service/WarnLogsCounter.java index b21f61fa88..2dff8e8e40 100644 --- a/spring-5-data-reactive/src/main/java/com/baeldung/tailablecursor/service/WarnLogsCounter.java +++ b/spring-5-data-reactive/src/main/java/com/baeldung/tailablecursor/service/WarnLogsCounter.java @@ -3,7 +3,7 @@ package com.baeldung.tailablecursor.service; import com.baeldung.tailablecursor.domain.Log; import com.baeldung.tailablecursor.domain.LogLevel; import lombok.extern.slf4j.Slf4j; -import org.springframework.data.mongodb.core.ReactiveMongoTemplate; +import org.springframework.data.mongodb.core.ReactiveMongoOperations; import reactor.core.Disposable; import reactor.core.publisher.Flux; @@ -21,10 +21,10 @@ public class WarnLogsCounter implements LogsCounter { private final AtomicInteger counter = new AtomicInteger(); private final Disposable subscription; - public WarnLogsCounter(ReactiveMongoTemplate template) { + public WarnLogsCounter(ReactiveMongoOperations template) { Flux stream = template.tail(query(where(LEVEL_FIELD_NAME).is(LogLevel.WARN)), Log.class); - subscription = stream.subscribe(l -> { - log.warn("WARN log received: " + l); + subscription = stream.subscribe(logEntity -> { + log.warn("WARN log received: " + logEntity); counter.incrementAndGet(); }); }