KaiSpace
tech

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有返回值。

处理任务执行关系:

  1. thenApply:相当于map
  2. thenCompose:相当于flatMap
  3. thenCombinea.thenCombine(b, Supplier)相当于combine
  4. thenRun(Runnable):执行runnable

前面的CompletableFuture标记为complete后再执行的操作

  1. runAfterBothcf1.runAfterBoth(cf2, c),cf1和cf2都标记为complete后再执行c。
  2. runAfterEithercf1.runAfterEither(cf2, c),cf1和cf2中一个标记为complete后就执行c。
  3. get():堵塞处理a.get()的线程直到a标记为complete,如果有exception会throw出来。
  4. 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.