Asynchronized programming in Java
多线程基础
我们经常听说四核八线程等,所谓线程如同前台一样,可以处理task。但是线程从哪里取register呢?register是由核心提供的,每个核心有自己的一组registers。
如果两个线程使用同一个核心,那么就需要registers上下文的切换,在现代优化中切换已经很快了,但是这样毕竟还是同时只能有一个线程在运行。这种设计并非不如单线程操作,因为我可以例如间断性地发送网络请求,或者进行磁盘读写操作,由于这些非常慢,单线程需要等待很久,这时我们切换到其他任务,这些与线程无关地操作就可以自己进行了。
如果使用不同核心,那么就是真正的并行计算了。
Java并行计算语法
基础
使用java.lang.Thread包
创建线程初始任务:new Thread(() -> {});
执行任务:.start()
new Thread(() -> {
for (int i = 1; i < 100; i += 1) {
System.out.print("_");
}
}).start();
new Thread(() -> {
for (int i = 2; i < 100; i += 1) {
System.out.print("*");
}
}).start();
暂停任务:Thread.sleep()为使当前运行这一命令的线程暂停
协调dependency
我们会使用Java.util.concurrent.CompletableFuture,这是一个Monad,存储了是否complete的side information。
新建一个任务:
CompletableFuture<Integer> a = CompletableFuture.completedFuture(taskA(x));新建一个已经执行完的操作。
CompletableFuture.runAsync(Runnable)新建后runAsync直接返回,但是runnable仍然在执行,当runnable执行结束之后将CompletableFuture mark as complete。
CompletableFuture.supplyAsync(Supplier)新建后supplyAsync直接返回,但是runnable仍然在执行,当runnable执行结束之后将CompletableFuture mark as complete。
CompletableFuture.runAsync(Runnable)和CompletableFuture.supplyAsync(Supplier)的区别:Runnable无返回值,Supplier有返回值。
处理任务执行关系:
thenApply:相当于mapthenCompose:相当于flatMapthenCombine:a.thenCombine(b, Supplier)相当于combinethenRun(Runnable):执行runnable
前面的CompletableFuture标记为complete后再执行的操作
runAfterBoth:cf1.runAfterBoth(cf2, c),cf1和cf2都标记为complete后再执行c。runAfterEither:cf1.runAfterEither(cf2, c),cf1和cf2中一个标记为complete后就执行c。get():堵塞处理a.get()的线程直到a标记为complete,如果有exception会throw出来。join():和get()一样,只是不throw exception。
并行计算例子:
CompletableFuture<Integer> a = CompletableFuture.completedFuture(taskA(x));
CompletableFuture<Integer> b = a.thenApply(x -> x + 1);
CompletableFuture<Integer> c = a.thenApply(x -> x + 1);
d = b.thenCombine(c, (x, y) -> x + y);
fork and join
对于recursive parallel execution,我们常常使用fork and join。我们先举一个例子:
import java.util.concurrent.RecursiveTask;
import java.util.concurrent.ForkJoinPool;
public class Fibonacci extends RecursiveTask<Integer> {
private final int n;
public Fibonacci(int n) {
this.n = n;
}
@Override
protected Integer compute() {
if (n <= 1) {
return n;
}
// 创建两个子任务
Fibonacci task1 = new Fibonacci(n - 1);
Fibonacci task2 = new Fibonacci(n - 2);
// 先让 task1 异步执行
task1.fork();
// 在当前线程中计算 task2
int result2 = task2.compute();
// 等待 task1 完成,并获取结果
int result1 = task1.join();
return result1 + result2;
}
public static void main(String[] args) {
int n = 10; // 要计算斐波那契数的项
ForkJoinPool pool = new ForkJoinPool();//创建task pool,这样才能加入task
Fibonacci fibTask = new Fibonacci(n);
int result = pool.invoke(fibTask); // 阻塞直至任务完成并返回结果
System.out.println("Fibonacci(" + n + ") = " + result);
}
}
recursiveTask有一个abstract method叫做compute(),指定了当fork and join pool处理这个任务时要进行什么操作。在这里,如果非tail call,则我们会创建两个任务fib(n - 1)和fib(n - 2)。
所谓fork操作其实只是给当前thread(处理fork的thread)的deque队首加入一个task。由于是recursive的,所以队首加入符合stack的性质。在这里我们加入了fib(n - 1)和fib(n - 2),实际上是先处理fib(n - 2)再处理fib(n - 1)。当然,如果我们的thread调用够快,可以在fib(n - 1)加入pool的一瞬间就处理,但实际上不可能这么快。所以由于fib(n - 2)在队首,所以先处理它。
join操作是堵塞当前thread(处理join的thread),并且给出返回值。
我们join的顺序和fork的顺序相反。
a.fork();
b.fork();
return b.join() + a.join();
这样的join顺序更efficient。如果我们先a.join(),那么无论哪一个线程(由于work stealing)在执行a,我们都要阻塞当前线程(处理a.join()的线程)到a执行结束。那么当前线程本来可能在处理b,现在被闲置了。而如果我们先b.join(),那么当前线程肯定是先处理b的,此时当前线程没有被闲置,而其他线程可以窃取当前线程的task deque的尾任务。
我们也可以将上面代码片段写成:
a.fork();
return b.compute() + a.join();
b.compute()是立刻进入线程内部处理,并阻塞线程,不允许其他task使用当前线程。
但是其实内部更加复杂,如果我们运行:
a.fork();
b.fork();
return a.join() + b.join();
b先在当前thread开始做,a被其他thread窃取开始工作。此时a.join()会堵塞当前线程,但是激活了help join的功能,它会帮助正在等待的线程工作,也就是从做a的线程中steal一些task,但是不会主动进行b的后续计算。而任务窃取需要开销,所以虽然两个线程都一直在工作且工作量一样,但是steal拖慢了计算,所以不如尽量各做各的。
而且如果在调用a.join()的时候a还没有被窃取,那么当前thread就会亲自首先执行join()任务。但是其中机制非常复杂,如果b还在被执行,那么可能还是会等a被别的thread stolen等等,但是当我们画出双端队列的图的时候会发现并不那么顺了,有些任务会滞留很久在队列中。
steal
其他thread会从当前thread的队尾取任务来做。这样设计是因为队尾的任务大概率工作量更大,如果每次steal小任务,则需要多次steal,而steal需要一定的时间,所以不如直接把大任务一同steal过来做完。
Thread和ForkJoinPool的区别
1. 线程管理模式不同
- 普通 Thread:
当你直接使用new Thread(runnable)创建线程时,你完全控制该线程的创建、启动以及终止。这个线程可以被指定用于特定的任务,不同的线程可以承担不同的责任。如果你需要组织这些线程的任务调度,则需要自己管理线程的状态、同步机制等。 - ForkJoinPool:
ForkJoinPool 是为分治任务(Divide and Conquer)的并行计算而设计的。它内部创建了一组专门的工作线程,这些线程会不断从任务队列中获取待处理的任务,并在必要时“窃取”(steal)其他线程的任务以提高利用率。这种机制能更高效地处理大量相互独立但又可以相互协作的任务。
2. 外部线程与 ForkJoinPool 内部线程的关系
- 外部线程不会自动成为 ForkJoinPool 的一部分:
如果你已经创建了一个普通的 Thread,比如用于执行一些特定逻辑,那么当你实例化一个 ForkJoinPool 时,这个外部线程不会被包含或纳入 ForkJoinPool 的内部工作线程中。ForkJoinPool 使用其自己内部的线程池来管理和执行任务。 - 调用线程的“帮助工作”作用:
有时,当你从一个非 ForkJoinPool 线程调用ForkJoinPool.commonPool().invoke(task)或者其他提交任务的方法时,这个调用线程会帮助完成部分任务。也就是说,虽然这个线程本身不是 ForkJoinPool 的一员,但在调用过程中,它可以参与到任务的执行中。不过,这并不代表它自动成为该池中的一员,而仅仅是临时参与工作。
3. 线程工厂和自定义线程
- ForkJoinPool 提供了自定义线程工厂的选项。你可以在创建 ForkJoinPool 时传入自定义的线程工厂,通过它来定制 ForkJoinWorkerThread 的创建方式(例如设置线程名称、优先级或其他属性)。但是,这里的定制仍然是针对 ForkJoinPool 内部创建的线程,而不是将外部已存在的线程“注入”到池中。
本段内容由chatgpt生成
Comments
No comments yet.