【问题标题】:Thread used for Java CompletableFuture composition?用于 Java CompletableFuture 组合的线程?
【发布时间】:2021-06-04 07:42:50
【问题描述】:

我开始对 Java CompletableFuture 组合感到满意,我曾使用过 JavaScript Promise。基本上,组合只是在指定的执行程序上安排了链式命令。但我不确定执行组合时哪个线程正在运行。

假设我有两个执行者,executor1executor2;为简单起见,假设它们是单独的线程池。我安排了CompletableFuture(使用非常宽松的描述):

CompletableFuture<Foo> futureFoo = CompletableFuture.supplyAsync(this::getFoo, executor1);

完成后,我使用第二个执行器将Foo 转换为Bar

CompletableFuture<Bar> futureBar .thenApplyAsync(this::fooToBar, executor2);

我了解getFoo() 将从executor1 线程池中的线程调用。我知道fooToBar() 将被executor2 线程池中的线程调用。

但是实际合成使用的是什么线程,即在getFoo()完成和futureFoo()完成之后;但是之前 fooToBar() 命令被安排在executor2 上? 换句话说,哪个线程实际运行代码以在第二个执行程序上安排第二个命令?

调度是否作为executor1 中调用getFoo() 的同一线程的一部分执行?如果是这样,这个可完成的未来组合是否等同于我在executor1 任务的第一个命令中手动安排fooToBar()

【问题讨论】:

  • 这是执行调用的线程,这取决于您的代码库。但是不,它既不是来自executor1 也不是executor2thenApplyAsync 的全部意义在于仅在这些操作中具有确定性。可能是more details here
  • "它是调用线程的线程......" 你到底是什么意思?你的意思是,如果在我的main() 方法中的一个名为main-thread 的线程中我调用原始的CompletableFuture.supplyAsync(),那么当它完成时,在executor2 中安排fooToBar() 操作的调用将在main-thread 中为好吧?但这怎么可能呢,因为main-thread 已经异步地走上了愉快的道路,现在正在分解素数或其他什么?
  • 我可能有点误解了你的问题,抱歉。您的问题是哪个线程将在executor2 上安排执行作为thenApplyAsync 的一部分?
  • 是的,@Eugene,这正是我的问题。 “换句话说,哪个线程实际运行代码以在第二个执行程序上调度第二个命令 [作为thenApplyAsync() 的一部分]?”
  • 我确信它只是完成了第一个未来的线程,我不明白为什么它应该/可能是任何其他的。随意阅读CompletableFutureCompletionStage 的来源。我刚刚做了,它看起来就像是第一个未来.. 代码读起来很糟糕,我无法确定它,所以我不想写答案。

标签: java parallel-processing completable-future executor


【解决方案1】:

这是故意未指定的。在实践中,当没有Async 后缀的变体被调用并表现出类似的行为时,它将由处理链式操作的相同代码处理。

所以当我们使用下面的测试代码时

CompletableFuture.supplyAsync(() -> {
    LockSupport.parkNanos(TimeUnit.SECONDS.toNanos(1));
    return "";
}, r -> new Thread(r, "A").start())
.thenAcceptAsync(s -> {}, r -> {
    System.out.println("scheduled by " + Thread.currentThread());
    new Thread(r, "B").start();
});

它可能会打印

scheduled by Thread[A,5,main]

因为完成前一阶段的线程用于调度依赖的操作。

但是当我们使用

CompletableFuture<String> first = CompletableFuture.supplyAsync(() -> "",
    r -> new Thread(r, "A").start());
LockSupport.parkNanos(TimeUnit.SECONDS.toNanos(1));
first.thenAcceptAsync(s -> {}, r -> {
    System.out.println("scheduled by " + Thread.currentThread());
    new Thread(r, "B").start();
});

它可能会打印

scheduled by Thread[main,5,main]

当主线程调用thenAcceptAsync 时,第一个future 已经完成,主线程将自行安排动作。

但这并不是故事的结局。当我们使用

CompletableFuture<String> first = CompletableFuture.supplyAsync(() -> {
    LockSupport.parkNanos(TimeUnit.MILLISECONDS.toNanos(5));
    return "";
}, r -> new Thread(r, "A").start());

Set<String> s = ConcurrentHashMap.newKeySet();
Runnable submitter = () -> {
    String n = Thread.currentThread().getName();
    do {
        for(int i = 0; i < 1000; i++)
            first.thenAcceptAsync(x -> s.add(n+" "+Thread.currentThread().getName()),
                Runnable::run);
    } while(!first.isDone());
};
Thread b = new Thread(submitter, "B");
Thread c = new Thread(submitter, "C");
b.start();
c.start();
b.join();
c.join();
System.out.println(s);

它可能不仅打印第一个场景中的B AC A 以及第二个场景中的B BC C 组合。在我的机器上,它还可以重现地打印B CC B 的组合,表明一个线程传递给thenAcceptAsync 的操作被另一个调用thenAcceptAsync 的线程同时使用不同的操作提交给了执行程序。

这与this answer 中描述的评估传递给thenApply(没有Async)的函数的线程场景相匹配。正如一开始所说,这是我所期望的,因为这两件事很可能由相同的代码处理。但与评估传递给thenApply 的函数的线程不同,在文档中甚至没有提到在Executor 上调用execute 方法的线程。所以理论上,另一个实现可以使用完全不同的线程,而不是调用未来的方法,也不会完成它。

【讨论】:

    【解决方案2】:

    最后是一个简单的程序,它确实喜欢您的代码 sn-p 并允许您使用它。

    当它等待的条件准备好时,输出确认您提供的执行程序被调用完成(除非您足够早地显式调用完成 - 这将发生在完成的调用线程中) - 上的 get() Future 阻塞,直到 Future 完成。

    提供一个 arg - 有一个 executor 1 和 executor 2,不提供 args 只有一个 executor。输出要么是(相同的执行器 - 事物在同一个执行器中作为单独的任务按顺序运行)-

    In thread Thread[main,5,main] - getFoo
    In thread Thread[main,5,main] - getFooToBar
    In thread Thread[pool-1-thread-1,5,main] - Supplying Foo
    In thread Thread[pool-1-thread-1,5,main] - fooToBar
    In thread Thread[main,5,main] - Completed
    

    OR(两个执行器 - 事情再次按顺序运行但使用不同的执行器)-

    In thread Thread[main,5,main] - getFoo
    In thread Thread[main,5,main] - getFooToBar
    In thread Thread[pool-1-thread-1,5,main] - Supplying Foo
    In thread Thread[pool-2-thread-1,5,main] - fooToBar
    In thread Thread[main,5,main] - Completed
    

    记住:带有执行程序的代码(在本例中可以立即在另一个线程中启动.. getFoo 在设置 FooToBar 之前就已被调用)。

    代码如下 -

    package your.test;
    
    import java.util.concurrent.CompletableFuture;
    import java.util.concurrent.ExecutionException;
    import java.util.concurrent.Executor;
    import java.util.concurrent.Executors;
    import java.util.concurrent.TimeUnit;
    import java.util.concurrent.TimeoutException;
    import java.util.function.Function;
    import java.util.function.Supplier;
    
    public class TestCompletableFuture {
        private static void dumpWhichThread(final String msg) {
            System.err.println("In thread " + Thread.currentThread().toString() + " - " + msg);
        }
    
        private static final class Foo {
            final int i;
            Foo(int i) {
                this.i = i;
            }
        };
        public static Supplier<Foo> getFoo() {
            dumpWhichThread("getFoo");
            return new Supplier<Foo>() {
                @Override
                public Foo get() {
                    dumpWhichThread("Supplying Foo");
                    return new Foo(10);
                }
    
            };
        }
    
        private static final class Bar {
            final String j;
            public Bar(final String j) {
                this.j = j;
            }
        };
        public static Function<Foo, Bar> getFooToBar() {
            dumpWhichThread("getFooToBar");
            return new Function<Foo, Bar>() {
                @Override
                public Bar apply(Foo t) {
                    dumpWhichThread("fooToBar");
                    return new Bar("" + t.i);
                }
            };
        }
    
    
        public static void main(final String args[]) throws InterruptedException, ExecutionException, TimeoutException {
            final TestCompletableFuture obj = new TestCompletableFuture();
            obj.running(args.length == 0);
        }
    
        private String running(final boolean sameExecutor) throws InterruptedException, ExecutionException, TimeoutException {
            final Executor executor1 = Executors.newSingleThreadExecutor(); 
            final Executor executor2 = sameExecutor ? executor1 : Executors.newSingleThreadExecutor(); 
            CompletableFuture<Foo> futureFoo = CompletableFuture.supplyAsync(getFoo(), executor1);
            CompletableFuture<Bar> futureBar = futureFoo.thenApplyAsync(getFooToBar(), executor2);
            try {
                // Try putting a complete here before the get ..
                return futureBar.get(50, TimeUnit.SECONDS).j;
            }
            finally {
                dumpWhichThread("Completed");
            }
        }
    }
    

    哪个线程触发 Bar 阶段进行 - 在上面 - 它是 executor1。一般来说,完成未来的线程(即给它一个值)是根据它释放事物的原因。如果您在主线程上立即完成 FutureFoo - 它会触发它。

    所以你必须小心这个。如果你有“N”个东西都在等待未来的结果——但只使用一个线程执行器——那么第一个调度的执行器将阻塞该执行器,直到它完成。您可以推断出 M 个线程、N 个未来 - 它可以衰减为“M”个锁,从而阻止其余的事情继续进行。

    【讨论】:

    • 您的示例显示了哪个线程执行这些功能,而不是哪个线程调度它们。
    • 它显示了 - getFoo 和 getFooToBar - 在主线程中完成调度时被调用(即在 running() 方法中 - 调用 supplyAsync 和 thenApplyAsync - 方法在调度时调用提供的)......确切的执行时间是所使用的特定 Executor 的函数 - 但绝对不会在主线程中发生......
    • 这与哪个线程调用supplyAsyncthenApplyAsync 无关。这是关于当第一个线程完成时哪个线程触发第二个依赖的未来。你错过了这个问题的重点。
    • 没有其他“神奇”线程... executor1 释放链接的执行,就像 executor2 所做的那样,如果 bar 然后链接到其他东西。
    • 一些有趣的事情(我之前没有使用过 CompletableFuture)我注意到如果你完成,而不是让异步完成,你确实得到了你期望传递给 futureBar 的东西,但是,执行器仍然运行 Supply 函数 .. 虽然它记录在 https://docs.oracle.com/javase/8/docs/api/java/util/concurrent/CompletableFuture.html 中,但它并不是我所期望的,Completeable 无法取消执行 - 所以是有道理的。
    猜你喜欢
    • 1970-01-01
    • 2021-02-20
    • 1970-01-01
    • 2016-10-28
    • 2018-12-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多