Модераторы: LSD, AntonSaburov
  

Поиск:

Ответ в темуСоздание новой темы Создание опроса
> Друзья, помогите с пулом потоков, Реализация пула потоков 
:(
    Опции темы
headzero
Дата 18.12.2008, 20:51 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Опытный
**


Профиль
Группа: Участник
Сообщений: 265
Регистрация: 14.11.2006

Репутация: нет
Всего: 1



Задача таковая - реализовать пул потоков, с заданным максимальным числом одновременно выполняющихся потоков + таск который выполняется и добавляет в очередь задач все новые и новые задачи.Работа пула не  должна зависеть от скорости выполняющихся тасков.

 Я пробовал что-то сделать. Вроде с горем пополам, но работает. Реализовал очень криво,иногда даже наугад, потому есть много вопросов:
1. В каких местах надо ставить синхронизацию?
2. Как считать количество выполняющихся на даный момент потоков? В каком месте нужно правильно инкрементировать/декрементировать  счетчик.
3. Когда задача мой таск выполняется слишком быстро, очередь потоков наполняется соответственно тоже быстро, но все время выполняется только один поток, а должно максимальное число.
4. Может надо использовать wait/notify для приостановки добавления задач если очередь задач пуста или достигнуто максимальное количество потоков?
Если можно, пожалуйста, помогите правильно реализовать, и подкорректируйте код. Спасибо
Код


public class TaskExecutor {
    private int thredPoolSize = 20; //Pool Size
    private static int activeThreadsnaumber;
    public Queue<Runnable> taskQueue = new LinkedList<Runnable>(); //Tasks to Execute
    public static int currentThreadsCounter;
    
    void begin() {
     while (true) {
         Runnable task = null;
        System.out.println("Queue size: "+taskQueue.size());
        
         activeThreadsnaumber = currentThreadsCounter;
         System.out.println("Vozmozhnoe chislo potokov: "+currentThreadsCounter);
         if (activeThreadsnaumber < thredPoolSize) {
                 if (taskQueue.size() != 0) {
                  synchronized (this) {
                            task = taskQueue.poll();
                            TeskExec thread = new TeskExec(this);
                            thread.setZadanie(task);
                            thread.start();
                           try {
                              Thread.sleep(10);
                            }
                            catch(InterruptedException e) {
                                
                            }
                           
                    }    
                
                }
       }
       }
  }
    public synchronized void addToQueue(Runnable task) {
       taskQueue.add(task);
     }
 
     public  void incrementNumOfThreads() {
      currentThreadsCounter++;
     }

     public  void decrementNumOfThreads() {
      currentThreadsCounter--;
     }

public static class Realozator implements Runnable { // Task to Run
    private String thName;
    public Thread t;
    private TaskExecutor ut;
    
    public Realozator(String thName,TaskExecutor ut) {
        this.thName = thName;
        t = new Thread(this,thName);
        this.ut = ut;
    }
    
    public void run() {
            for (int i=0;i<10;i++) {
                try {
                        this.t.sleep(10);
                   }
                   catch(InterruptedException e) {
                                   //TODO        
                   }
                   ut.addToQueue(new Realozator("Inside "+currentThreadsCounter,ut));
                   System.out.println("This is thread <"+thName+"> vipolnen");
            }
        
       
    }

}


public static class TeskExec extends Thread {
    private Runnable task;
    private TaskExecutor inut;
    private Thread th;
    public TeskExec(TaskExecutor inut) {
        this.inut = inut;
    }
    public synchronized void setZadanie(Runnable task) {
        this.task = task;
        th = new Thread(task);
    }
    
    
    public void run() {  
        inut.incrementNumOfThreads();
          th.run();
             synchronized(this) {
             task = null;
              try {
                  th.join();
            }  catch (InterruptedException e) {
                  e.printStackTrace();
            }
          }
             inut.decrementNumOfThreads();
             //inut.notify();
             
    }
}

  public static void main(String[] args) {
      TaskExecutor ut = new TaskExecutor();
      ut.addToQueue(new Realozator("My",ut));
      ut.begin();
    
    }

}









--------------------
Воображение важнее знания
                                                     (Алберт Эйнштейн)
PM MAIL   Вверх
Vurn
Дата 18.12.2008, 22:25 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Шустрый
*


Профиль
Группа: Участник
Сообщений: 51
Регистрация: 24.5.2007

Репутация: 2
Всего: 3



PM MAIL   Вверх
Бонифаций
Дата 19.12.2008, 02:11 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Опытный
**


Профиль
Группа: Участник
Сообщений: 827
Регистрация: 15.9.2005
Где: Brisbane

Репутация: 1
Всего: 40



это есть в стандартной бибке явы. Смотрите Executors.

или вот вам ссылка

http://java.sun.com/javase/6/docs/api/java...ThreadPool(int)



--------------------
 Бонифаций.
 
PM MAIL ICQ Skype GTalk Jabber YIM   Вверх
headzero
Дата 19.12.2008, 10:17 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Опытный
**


Профиль
Группа: Участник
Сообщений: 265
Регистрация: 14.11.2006

Репутация: нет
Всего: 1



Цитата

Смотрите Executors


Спасибо, про эти пакеты я знаю. Но задача стоит в том что-бы реализовать пул без сторонних пакетов, используя только примитивы синхронизации и стандартные функции из Thread.


--------------------
Воображение важнее знания
                                                     (Алберт Эйнштейн)
PM MAIL   Вверх
SoulKeeper
Дата 19.12.2008, 11:47 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Опытный
**


Профиль
Группа: Участник
Сообщений: 375
Регистрация: 14.1.2007
Где: Ukraine, Lviv.

Репутация: 11
Всего: 15



Этот пул - стандартные пакеты джавы.
PM MAIL   Вверх
headzero
Дата 19.12.2008, 15:16 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Опытный
**


Профиль
Группа: Участник
Сообщений: 265
Регистрация: 14.11.2006

Репутация: нет
Всего: 1



Задача состоит в том чтобы реализовать пул самостоятельно, без пакета concurrency.Жду предложений, и комментариев.


--------------------
Воображение важнее знания
                                                     (Алберт Эйнштейн)
PM MAIL   Вверх
Vurn
Дата 19.12.2008, 16:06 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Шустрый
*


Профиль
Группа: Участник
Сообщений: 51
Регистрация: 24.5.2007

Репутация: 2
Всего: 3



А что тут сложного? только очередь переделать на Queue, обработчики будут спать на ней, вычислители, захватив - будут будить, обработчики обработав кусок - уменьшать переменную обратного отсчета и будить вычислителей.
Код

import java.util.Formatter;
import java.util.Queue;
import java.util.LinkedList;


/**
 * Главный класс - Обработка одного задания
 */
public class ParallelComputation2 implements Runnable{
    final Queue<CompPart> query;  // очередь запросов на обработку
    int countDown = 0;
    private static final int PROCESSORS = Runtime.getRuntime().availableProcessors(); //константа-количество процессоров

    public ParallelComputation2(Queue<CompPart> query) {
        this.query = query;
    }

    public void run() {
        final int taskNo = Integer.parseInt(Thread.currentThread().getName());
        double[][] data = new double[10][10];// данные
        for (int i = 0; i < data.length; i++ ) {
            double[] row = data[i];
            for (int j = 0; j < row.length; j++) {
                row[j] = i + j + taskNo;
            }
        }
        synchronized (query) { //захват очереди
            countDown = data.length; // переменная обратного отсчета 
            for (int i = 0; i < data.length; i++) {
                // просто пихаем в очередь часть, обработчики выдерут по очереди и обработают
                query.add(new CompPart(i, data, this));
            }
            query.notifyAll(); // будим обработчиков (если дрыхнут в ожидании данных)
        }
        synchronized (this) { // захват нашего объекта
            // может быть ситуация, когда запустив обработчиков, планировщик так долго не давал
            // этой нити отработать, что обработчики все сделали. Данная проверка исключает ошибочное ожидание на 
            //этом объекте в такой ситуации.
            while (countDown != 0) {  // while - потому что JVM может ошибочно разбудить поток раньше времени.
                try {
                    this.wait(); // ну, а если хотя бы одна часть необработана - дрыхнем, пока нас не разбудят.
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        }
        // выводная форма собирается через StringBuilder внутри класса Formatter, чтобы не смешались
        // данные на печати из разных нитей
        Formatter f = new Formatter();
        f.format("All parts processed, show result thread ").format(Thread.currentThread().getName()).format(" \n");
        for (double[] row :data) {
            for (double aRow : row) {
                f.format("%4.1f ",aRow);
            }
            f.format("\n");
        }
        System.out.println(f); // один-единственный синхронный вывод на печать.

        // извлекаем ответ и обрабатываем
    }

    public static void main(String[] args) {
        // создаем единую нить запросов
        final Queue<CompPart> commonQuery = new LinkedList<CompPart>();
        // создаем обработчики по числу процессоров на компьютере
        for (int i = 0; i < PROCESSORS; i++) {
            Thread th = new Thread(new Processor2(commonQuery)); //создаем нить обработчика
            th.setDaemon(true);
            th.start();
        }

        // создаем 5-ть параллельных задач на обработку, разделяющих одну очередь и одних и тех же обработчиков
        for (int i = 0; i < 5; i++) {
            new Thread(new ParallelComputation2(commonQuery), String.valueOf(i)).start();
        }
    }
    /**
     * класс-обработчик
     */
    static class Processor2 implements Runnable{
        final Queue<CompPart> income; // входная очередь

        /**
         * конструктор
         * @param income
         */
        Processor2(Queue<CompPart> income) {
            this.income = income;
        }

        /**
         * основной метод ожидания/обработки
         */
        public void run() {
            try {
                // бесконечный цикл можно прервать поставив нить данного обработчика в .interrupt()
                // но так как она daemon, можно забить и просто выйти из программы.
                while(!Thread.interrupted()) {
                    CompPart cp = null; // подготавливаем объект для извлечения запроса на обработку
                    synchronized (income) { // захват очерди
                        // Если очередь пустая - спим на ней, пока вычислители не накидают данные и не разбудят
                        if (income.isEmpty()) { 
                            income.wait();
                        }
                        cp = income.poll(); // извлекаем объект
                    }
                    // может получиться так, что был кинут один объект, а на очереди спало обработчиков
                    // объект будет захвачен одним обработчиком, остальные получат null и проскочат цикл
                    if (cp != null) {
                        double[] matrixRow = cp.getDataRef()[cp.getRow()]; // обрабатываемый ряд
                        {
                        }
                        // process....
                        for (int i = 0; i < matrixRow.length; i++) {
                            matrixRow[i] += 2;
                        }

                        // inform вычислителя
                        final ParallelComputation2 pc2;
                        pc2 = cp.getPc2();
                        synchronized (pc2) { // захват вычислителя
                            pc2.countDown--; // уменьшаем переменную обратного отсчета
                            if (pc2.countDown == 0) { // если больше от вычислителя задач не будет...
                                pc2.notify(); // будим его (если он спит)
                            }
                        }
                    }
                }
            } catch (InterruptedException ignored) {}

        }
    }

    /**
     * дополнительный класс-сообщение для передачи по очереди.
     */
    static class CompPart {
        private int row;
        private double[][] dataRef;
        final private ParallelComputation2 pc2;

        CompPart(int row, double[][] dataRef, ParallelComputation2 pc2) {
            this.row = row;
            this.dataRef = dataRef;
            this.pc2 = pc2;
        }

        public int getRow() {
            return row;
        }

        public double[][] getDataRef() {
            return dataRef;
        }

        public ParallelComputation2 getPc2() {
            return pc2;
        }
    }
}


Это сообщение отредактировал(а) Vurn - 19.12.2008, 16:20
PM MAIL   Вверх
Vurn
Дата 19.12.2008, 16:54 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Шустрый
*


Профиль
Группа: Участник
Сообщений: 51
Регистрация: 24.5.2007

Репутация: 2
Всего: 3



Кстати, никто не мешал в той ссылке, что я давал выше, написать свои CountDownLatch & BlockingQueue, они пишутся элементарно чрезе synchronized/wait/notify.
P.S. от нефиг делать написал.
Код

public class BlockingQueue<T> {
    private Entry<T> head = null, tail = null;

    public synchronized  void put(T val) throws InterruptedException {
        Entry<T> e = new Entry<T>();
        e.value = val;
        if (head == null) {
            head = tail = e;
        } else {
            tail.next = e;
            tail = e;
        }
        notify();
    }
    public synchronized T take() throws InterruptedException {
        while (true) {
            if (head != null) {
                Entry<T> temp = head;
                if (head == tail) {
                    head = tail = null;
                } else {
                    head = head.next;
                }
                return temp.value;
            } else {
                wait();
            }
        }
    }
    private class Entry<T> {
        private Entry<T> next;
        private T value;
    }
}




Код

public class CountDownLatch {
    private int countDown;

    public CountDownLatch(int countDown) {
        this.countDown = countDown;
    }
    public synchronized void await() throws InterruptedException {
        while (countDown != 0) {
            wait();
        }
    }
    public synchronized void countDown() {
        if (countDown > 0) {
            countDown--;
        }
        if (countDown == 0) {
            notifyAll();
        } 

    }
}



Это сообщение отредактировал(а) Vurn - 19.12.2008, 18:04
PM MAIL   Вверх
  
Ответ в темуСоздание новой темы Создание опроса
Правила форума "Java"
LSD   AntonSaburov
powerOn   tux
javastic
  • Прежде, чем задать вопрос, прочтите это!
  • Книги по Java собираются здесь.
  • Документация и ресурсы по Java находятся здесь.
  • Используйте теги [code=java][/code] для подсветки кода. Используйтe чекбокс "транслит", если у Вас нет русских шрифтов.
  • Помечайте свой вопрос как решённый, если на него получен ответ. Ссылка "Пометить как решённый" находится над первым постом.
  • Действия модераторов можно обсудить здесь.
  • FAQ раздела лежит здесь.

Если Вам помогли, и атмосфера форума Вам понравилась, то заходите к нам чаще! С уважением, LSD, AntonSaburov, powerOn, tux, javastic.

 
0 Пользователей читают эту тему (0 Гостей и 0 Скрытых Пользователей)
0 Пользователей:
« Предыдущая тема | Java: Общие вопросы | Следующая тема »


 




[ Время генерации скрипта: 0.0627 ]   [ Использовано запросов: 22 ]   [ GZIP включён ]


Реклама на сайте     Информационное спонсорство

 
По вопросам размещения рекламы пишите на vladimir(sobaka)vingrad.ru
Отказ от ответственности     Powered by Invision Power Board(R) 1.3 © 2003  IPS, Inc.