我的解决方案主题是(它可以与 JDK 9+ 一起使用,因为自该版本以来公开了几个可覆盖的方法)
让整个生态系统了解 MDC
为此,我们需要解决以下场景:
-
什么时候我们可以从这个类中获得 CompletableFuture 的新实例? → 我们需要返回一个 MDC 感知版本。
-
什么时候我们才能从这个类之外获得 CompletableFuture 的新实例? → 我们需要返回相同的 MDC 感知版本。
-
在 CompletableFuture 类中使用哪个执行器? → 在所有情况下,我们都需要确保所有执行器都支持 MDC
为此,让我们通过扩展它来创建CompletableFuture 的MDC 感知版本类。我的版本如下所示
import org.slf4j.MDC;
import java.util.Map;
import java.util.concurrent.*;
import java.util.function.Function;
import java.util.function.Supplier;
public class MDCAwareCompletableFuture<T> extends CompletableFuture<T> {
public static final ExecutorService MDC_AWARE_ASYNC_POOL = new MDCAwareForkJoinPool();
@Override
public CompletableFuture newIncompleteFuture() {
return new MDCAwareCompletableFuture();
}
@Override
public Executor defaultExecutor() {
return MDC_AWARE_ASYNC_POOL;
}
public static <T> CompletionStage<T> getMDCAwareCompletionStage(CompletableFuture<T> future) {
return new MDCAwareCompletableFuture<>()
.completeAsync(() -> null)
.thenCombineAsync(future, (aVoid, value) -> value);
}
public static <T> CompletionStage<T> getMDCHandledCompletionStage(CompletableFuture<T> future,
Function<Throwable, T> throwableFunction) {
Map<String, String> contextMap = MDC.getCopyOfContextMap();
return getMDCAwareCompletionStage(future)
.handle((value, throwable) -> {
setMDCContext(contextMap);
if (throwable != null) {
return throwableFunction.apply(throwable);
}
return value;
});
}
}
MDCAwareForkJoinPool 类看起来像(为简单起见,跳过了带有 ForkJoinTask 参数的方法)
public class MDCAwareForkJoinPool extends ForkJoinPool {
//Override constructors which you need
@Override
public <T> ForkJoinTask<T> submit(Callable<T> task) {
return super.submit(MDCUtility.wrapWithMdcContext(task));
}
@Override
public <T> ForkJoinTask<T> submit(Runnable task, T result) {
return super.submit(wrapWithMdcContext(task), result);
}
@Override
public ForkJoinTask<?> submit(Runnable task) {
return super.submit(wrapWithMdcContext(task));
}
@Override
public void execute(Runnable task) {
super.execute(wrapWithMdcContext(task));
}
}
包装的实用方法是这样的
public static <T> Callable<T> wrapWithMdcContext(Callable<T> task) {
//save the current MDC context
Map<String, String> contextMap = MDC.getCopyOfContextMap();
return () -> {
setMDCContext(contextMap);
try {
return task.call();
} finally {
// once the task is complete, clear MDC
MDC.clear();
}
};
}
public static Runnable wrapWithMdcContext(Runnable task) {
//save the current MDC context
Map<String, String> contextMap = MDC.getCopyOfContextMap();
return () -> {
setMDCContext(contextMap);
try {
return task.run();
} finally {
// once the task is complete, clear MDC
MDC.clear();
}
};
}
public static void setMDCContext(Map<String, String> contextMap) {
MDC.clear();
if (contextMap != null) {
MDC.setContextMap(contextMap);
}
}
以下是一些使用指南:
- 使用
MDCAwareCompletableFuture 类而不是CompletableFuture 类。
-
CompletableFuture 类中的几个方法实例化了 self 版本,例如 new CompletableFuture...。对于此类方法(大多数公共静态方法),使用替代方法获取MDCAwareCompletableFuture 的实例。使用替代方法的示例可能不是使用CompletableFuture.supplyAsync(...),您可以选择new MDCAwareCompletableFuture<>().completeAsync(...)
- 当您因为某个外部库返回
CompletableFuture 的实例而遇到困难时,请使用getMDCAwareCompletionStage 方法将CompletableFuture 的实例转换为MDCAwareCompletableFuture。显然,您不能在该库中保留上下文,但在您的代码命中应用程序代码后,此方法仍会保留上下文。
- 在提供执行程序作为参数时,请确保它是 MDC 感知的,例如
MDCAwareForkJoinPool。您也可以通过覆盖 execute 方法来创建 MDCAwareThreadPoolExecutor 以服务于您的用例。你明白了!
这样,您的代码将如下所示
List<CompletableFuture<UpdateHotelAllotmentsRsp>> futures =
tasks.stream()
new MDCAwareCompletableFuture<UpdateHotelAllotmentsRsp>().completeAsync(
() -> businesslogic(task))
.collect(Collectors.toList());
List results = futures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList());
public UpdateHotelAllotmentsRsp businesslogic(Task task) {
LOGGER.info("mdc fishtag context is not lost here");
}
您可以在post 中找到上述所有内容的详细解释。