diff --git a/src/main/java/io/ebean/BackgroundExecutor.java b/src/main/java/io/ebean/BackgroundExecutor.java index b88f8b658..f6439773e 100644 --- a/src/main/java/io/ebean/BackgroundExecutor.java +++ b/src/main/java/io/ebean/BackgroundExecutor.java @@ -1,6 +1,8 @@ package io.ebean; +import java.util.concurrent.Callable; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; /** @@ -37,4 +39,21 @@ public interface BackgroundExecutor { *

*/ void executePeriodically(Runnable r, long delay, TimeUnit unit); + + /** + * Schedules a Runnable for one-shot action that becomes enabled after the given delay. + * + * @return a ScheduledFuture representing pending completion of the task and + * whose get() method will return null upon completion + */ + ScheduledFuture schedule(Runnable r, long delay, TimeUnit unit); + + /** + * Schedules a Callable for one-shot action that becomes enabled after the given delay. + * + * @return a ScheduledFuture that can be used to extract result or cancel + */ + ScheduledFuture schedule(Callable c, long delay, TimeUnit unit); + + } diff --git a/src/main/java/io/ebeaninternal/server/core/DefaultBackgroundExecutor.java b/src/main/java/io/ebeaninternal/server/core/DefaultBackgroundExecutor.java index 200f7fedd..e458c5165 100644 --- a/src/main/java/io/ebeaninternal/server/core/DefaultBackgroundExecutor.java +++ b/src/main/java/io/ebeaninternal/server/core/DefaultBackgroundExecutor.java @@ -6,6 +6,8 @@ import io.ebeaninternal.server.lib.DaemonScheduleThreadPool; import org.slf4j.MDC; import java.util.Map; +import java.util.concurrent.Callable; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; /** @@ -64,6 +66,46 @@ public class DefaultBackgroundExecutor implements SpiBackgroundExecutor { } } + @Override + public ScheduledFuture schedule(Runnable r, long delay, TimeUnit unit) { + final Map map = MDC.getCopyOfContextMap(); + + if (map == null) { + return schedulePool.schedule(r, delay, unit); + } else { + return schedulePool.schedule(() -> { + MDC.setContextMap(map); + try { + r.run(); + } finally { + MDC.clear(); + } + }, delay, unit); + } + } + + @Override + public ScheduledFuture schedule(Callable c, long delay, TimeUnit unit) { + final Map map = MDC.getCopyOfContextMap(); + + if (map == null) { + return schedulePool.schedule(c, delay, unit); + } else { + return schedulePool.schedule(new Callable() { + @Override + public V call() throws Exception { + MDC.setContextMap(map); + try { + return c.call(); + } finally { + MDC.clear(); + } + } + + }, delay, unit); + } + } + @Override public void shutdown() { pool.shutdown();