【问题标题】:Java multiple thread join issueJava多线程连接问题
【发布时间】:2023-03-26 11:06:01
【问题描述】:

所以我需要使用线程(已经拆分)处理几个数据文件,并且我在如何停止主线程直到所有子线程完成时遇到问题。 我环顾四周并尝试使用 join() 但这会导致问题:

  • 如果我加入主线程和最后一个线程,那么由于其他线程同时运行,最后一个线程并不总是最后一个完成
  • 如果我将主线程与所有其他线程一起加入,那么它们不会同时运行,第二个需要第一个先完成。 还尝试了 wait() 和 notify() 但有更多问题。这是我的代码的一部分

        public class Matrix extends MapReduce {
        ArrayList<String> VecteurLines = new ArrayList<String>();
        protected int[] nbrLnCol = {0,0};
        protected static double[] res;

        public Matrix(String n) {
            super(n);
        }
        public Matrix(String n,String m){
            super(n,m);
        }
    public void Reduce() throws IOException, InterruptedException, MatrixException {

            for (int i = 1; i <= Chunks; i++) {

                Thread t=new Thread(new RunThread(VecteurLines,i,this));
                t.start();

            }
        }

这是处理线程的类


    public class RunThread extends Matrix implements Runnable {
            Matrix ma;
            ArrayList<String> vec;
            int threadNbr;


            public RunThread(ArrayList<String> vec, int threadNbr,Matrix ma)  {
                super("","");
                this.vec=vec;this.threadNbr=threadNbr;this.ma=ma; }

            @Override
            public void run() {

                FileInputStream fin = null;
                try {
                    fin = new FileInputStream(ma.getNom()+threadNbr+".txt");
                } catch (FileNotFoundException e) {
                    e.printStackTrace();
                }
                Scanner sc = new Scanner(fin);


                while (sc.hasNext()) {
                    String nextString = sc.next();

                    ma.nbrLnCol[0]++;
                    String [] arr = nextString.split(",");
                    ma.nbrLnCol[1]=arr.length;
                    double c=0;
                    for(int j=0;j<arr.length;j++)
                    {
                        c+=(Double.parseDouble(arr[j])*Double.parseDouble(vec.get(j)));

                    }

                    res[threadNbr-1]=c;
                }
                sc.close();
                try {
                    fin.close();
                } catch (IOException e) {
                    e.printStackTrace();
                }

                File file = new File(ma.getNom()+threadNbr+".txt");
                file.delete();
            }

【问题讨论】:

  • 显而易见的解决方案是加入所有子线程,但我在您的代码中看不到 join()。
  • 我删除了它,因为它最终会一个接一个地运行线程,而不是同时运行所有线程
  • join 不应导致任何线程延迟。你到底尝试了什么?主线程中发生了什么?看起来RunThread 在每个线程中都执行了一些操作,但是您提供的代码除了启动一系列线程之外,没有任何线程管理。
  • @Ayman Elya join() 不会让线程一一运行。有不同的原因。
  • 那段代码,有一个错误,但在我们得到代码之前我们无法修复它。

标签: java multithreading java-threads


【解决方案1】:

试试这样:

 private List<Thread> threadList = new ArrayList<>();

 public void Reduce() {
     threadList.clear();
     for (int i = 1; i <= Chunks; i++) {
         Thread t  =new Thread(new RunThread(VecteurLines,i,this));
         threadList.add(t);
     }

     // start all worker threads
     for(int i=0; i<threadList.size(); i++){
         threadList.get(i).start();
     }

     // wait until all worker threads is finished
     while (true) {
         int threadIsNotLive = 0;
         for (int i = 0; i < threadList.size(); i++) {
             Thread t = threadList.get(i);
             if (!t.isAlive() || t == null) {
                 ++threadIsNotLive;
             }
         }
         if(threadIsNotLive>0 && (threadList.size() == threadIsNotLive)){
             break;
             // all worker threads is finished
         }
         else {
             Thread.sleep(50);
             // wait until all worker threads is finished
         }
     }
 }

或

 public void Reduce() {
     List<Thread> threadList = new ArrayList<>();
     for (int i = 1; i <= Chunks; i++) {
         Thread t  =new Thread(new RunThread(VecteurLines,i,this));
         threadList.add(t);
     }

     // start all worker threads
     for(int i=0; i<threadList.size(); i++){
         threadList.get(i).start();
         threadList.get(i).join();
     }
 }

【讨论】:

  • 将所有Thread 实例存储到一个列表中是正确的方法,以确保在主线程开始等待之前所有实例都已启动。然而,做那个复杂的最后一个循环是没有意义的。一个简单的for(Thread t: threadList) t.join(); 就可以了。而不是将列表存储到字段中,需要clear() 它,只需将List&lt;Thread&gt; threadList = new ArrayList&lt;&gt;(); 声明为方法中的第一条语句,使其成为局部变量。
  • 不,使用该代码,您重新引入了问题,在启动后立即等待每个线程使得执行串行。您需要一个循环来启动所有线程,然后需要另一个循环来等待它们。并且不要犹豫使用 for-each 循环 for(Thread t: threadList) … 而不是索引循环。
【解决方案2】:

我相信您的代码需要两点: 你的主线程必须在所有线程执行完毕后最后结束,因为你说

“如何停止主线程直到所有子线程完成”

。 其次,线程应该一个接一个地完成,即第二个线程应该在第一个线程之后完成,就像你说的那样

“第二个需要第一个先完成。”

这是我使用 join 执行此操作的代码。

public class Matrix extends MapReduce {
    ArrayList<String> VecteurLines = new ArrayList<String>();
    protected int[] nbrLnCol = {0,0};
    protected static double[] res;

    public Matrix(String n) {
        super(n);
    }
    public Matrix(String n,String m){
        super(n,m);
    }
public void Reduce() throws IOException, InterruptedException, MatrixException {
    Thread t = null;
        for (int i = 1; i <= Chunks; i++) {

            Thread t=new Thread(new RunThread(t,VecteurLines,i,this));
            t.start();

        }
      t.join(); // finally main thread joining with the last thread.
    }

和

public class RunThread extends Matrix implements Runnable {
        Matrix ma;
        ArrayList<String> vec;
        int threadNbr;
        Thread t;


        public RunThread(t,ArrayList<String> vec, int threadNbr,Matrix ma)  {
            this.t = t;
            super("","");
            this.vec=vec;this.threadNbr=threadNbr;this.ma=ma; }

        @Override
        public void run() {                
            FileInputStream fin = null;
            try {
                fin = new FileInputStream(ma.getNom()+threadNbr+".txt");
            } catch (FileNotFoundException e) {
                e.printStackTrace();
            }
            Scanner sc = new Scanner(fin);


            while (sc.hasNext()) {
                String nextString = sc.next();

                ma.nbrLnCol[0]++;
                String [] arr = nextString.split(",");
                ma.nbrLnCol[1]=arr.length;
                double c=0;
                for(int j=0;j<arr.length;j++)
                {
                    c+=(Double.parseDouble(arr[j])*Double.parseDouble(vec.get(j)));

                }

                res[threadNbr-1]=c;
            }
            sc.close();
            try {
                fin.close();
            } catch (IOException e) {
                e.printStackTrace();
            }

            File file = new File(ma.getNom()+threadNbr+".txt");
            file.delete();
            if(t!=null){
             t.join(); //join with the previous thread eg. thread2 joining with thread1
            }
        }

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-11-14
    • 2020-09-17
    • 1970-01-01
    • 1970-01-01
    • 2013-03-06
    • 2016-06-23
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多