适配Guava的AbstractExecutorService实现gRPC上下文传播的最优方案探讨
适配Guava的AbstractExecutorService实现gRPC上下文传播的最优方案探讨
嘿,我太懂这种找不到现成轮子只能自己造的痛苦了!当初我想在gRPC Context库里找个能直接传播上下文的ExecutorService实现,结果也是啥都没找到,只能自己动手封装一个,看到你的思路简直跟我当时一模一样,咱们来好好唠唠最优解~
首先先明确需求:我们要把任意ExecutorService包装一层,让它在提交任务时,自动把提交线程的当前gRPC上下文绑定到任务上,这样任务执行时就能复用这个上下文了。
最开始你想到的两个方案各有坑:
- 只在
newTaskFor()里做包装:这方法看着挺规整,毕竟newTaskFor()就是用来创建任务的入口,但问题是如果有人直接调用execute()方法提交Runnable,这个路径就绕开了newTaskFor(),任务没被包装,上下文直接丢了,完全达不到目的。 - 只在
execute()里做包装:这倒是能覆盖直接调用execute()的场景,但你担心的点特别对——万一有个嵌套的ExecutorService在工作线程里调用execute(),那捕获的就是工作线程的上下文,根本不是我们要的提交线程的上下文,直接搞反了。
不过你后来想到的那个带自定义任务+类型检测的方案,绝对是最优解!我来给你拆解下为啥这个方案靠谱:
核心思路
- 自定义任务类:写一个
GrpcPropagatingTask继承FutureTask,在构造函数里直接用当前线程的gRPC上下文把传入的Callable或Runnable包起来。这样只要是通过newTaskFor()创建的任务,天生就带着提交线程的上下文。 - 重写
newTaskFor():让它返回我们的自定义任务,这样所有通过submit()、invokeAll()等方法提交的任务,都会走这个路径被正确包装。 - 增强
execute()方法:在execute()里先判断传入的Runnable是不是已经是我们的自定义任务——如果是,说明已经被包装过了,直接交给底层ExecutorService执行;如果不是,就用当前线程的上下文包装后再提交,完美覆盖直接调用execute()的场景。
完整实现代码
import io.grpc.Context; import java.util.concurrent.AbstractExecutorService; import java.util.concurrent.Callable; import java.util.concurrent.ExecutorService; import java.util.concurrent.FutureTask; import static com.google.common.base.Preconditions.checkNotNull; public final class GrpcPropagatingExecutorService extends AbstractExecutorService { private final ExecutorService delegate; public GrpcPropagatingExecutorService(ExecutorService delegate) { this.delegate = checkNotNull(delegate); } // 自定义带上下文传播的任务类 private static final class GrpcPropagatingTask<V> extends FutureTask<V> { public GrpcPropagatingTask(Callable<V> callable) { super(Context.current().wrap(callable)); } public GrpcPropagatingTask(Runnable runnable, V result) { super(Context.current().wrap(runnable), result); } } @Override protected <T> RunnableFuture<T> newTaskFor(Runnable runnable, T value) { return new GrpcPropagatingTask<>(runnable, value); } @Override protected <T> RunnableFuture<T> newTaskFor(Callable<T> callable) { return new GrpcPropagatingTask<>(callable); } @Override public void execute(Runnable command) { // 避免重复包装,只处理非自定义任务 if (!(command instanceof GrpcPropagatingTask)) { command = Context.current().wrap(command); } delegate.execute(command); } // 其他AbstractExecutorService要求实现的方法直接委托给delegate即可 @Override public void shutdown() { delegate.shutdown(); } @Override public boolean isShutdown() { return delegate.isShutdown(); } @Override public boolean isTerminated() { return delegate.isTerminated(); } @Override public boolean awaitTermination(long timeout, java.util.concurrent.TimeUnit unit) throws InterruptedException { return delegate.awaitTermination(timeout, unit); } @Override public java.util.List<Runnable> shutdownNow() { return delegate.shutdownNow(); } }
方案优势
- 全覆盖无遗漏:不管是通过
submit()系列方法还是直接调用execute()提交任务,都能被正确包装,不会漏掉任何场景。 - 无重复包装:通过类型检测避免了同一任务被多次包装的问题,性能上没额外负担。
- 上下文正确性:不管是
newTaskFor()还是execute(),都是在提交任务的线程中被调用的(按照ExecutorService的规范,工作线程不会用来提交任务),所以捕获的肯定是我们需要的提交线程的gRPC上下文,不会搞混。
总结下来,这个方案就是“双保险”又不冗余的最优解,完全解决了你之前担心的所有问题~
备注:内容来源于stack exchange,提问作者David
相关产品推荐
相关产品推荐

