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

Поиск:

Ответ в темуСоздание новой темы Создание опроса
> Поиск в ширину на потоках, мультитред версия 
V
    Опции темы
Platon
Дата 25.9.2008, 06:42 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Эксперт
***


Профиль
Группа: Завсегдатай
Сообщений: 1801
Регистрация: 25.4.2006

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



Здравствуйте, уважаемые.

Возникла потребность в написании многопоточного алгоритма поиска в ширину.

Код

        LinkedList<Data> datas = new LinkedList<Data>();
        datas.add(url);
        while (!datas.isEmpty()) {
            Data cur = datas.removeFirst();
            processData(cur, datas);
        }

private static void processData(Data cur, LinkedList<Data> datas) {
    for (/*something related to cur data*/) {
        if (/* item in something related to cur data is good enough */)
            datas.add(/* item in something related to cur data */);
    }
}


Моя недоделанная версия:

Код

        DataQueue queue = new DataQueue();

        Thread[] threads = new Thread[threadsAmount];
        for (int i = 0; i < threadsAmount; i++)
            threads[i] = new Thread(new Runnable() {
                public void run() {
                    Thread t = Thread.currentThread();
                    try {
                        while (!t.isInterrupted()) {
                            Data cur = queue.get();
                            processData(cur, queue);
                            queue.dataProcessed(cur);
                        }
                    } catch (InterruptedException e) {

                    }
                }
            });


DataQueue хранит в себе ArrayBlockingQueue get/add соответствуют take/add

Несложное решение, но не хватает одного - как прервать одновременно все потоки по завершение обработки всех объектов данных? В моем коде на это гипотетически способен только DataQueue, но как сделать, я не знаю.

Добавлено через 13 минут и 22 секунды
Цитата(Platon @  25.9.2008,  07:42 Найти цитируемый пост)
Возникла потребность в написании

как сказал, прямо программистский наркоман!!!

Это сообщение отредактировал(а) Platon - 25.9.2008, 06:43
PM MAIL ICQ   Вверх
Platon
Дата 25.9.2008, 10:55 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Эксперт
***


Профиль
Группа: Завсегдатай
Сообщений: 1801
Регистрация: 25.4.2006

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



Сам заказал - сам спляшу:

Код

import java.util.concurrent.ArrayBlockingQueue;

public class ThreadedBFSTest {
    public static void main(String[] args) throws InterruptedException {
        Thread[] t = new Thread[100];
        Producer p = new Producer();
        p.put(new Object());
        for (int i = 0; i < t.length; i++)
            t[i] = new Thread(new Consumer(p));
        for (Thread tt : t)
            tt.start();
        for (Thread tt : t)
            tt.join();
        System.out.println("Done");
    }

    private static class Producer {

        private ArrayBlockingQueue<Object> queue = new ArrayBlockingQueue<Object>(4096);

        private int queued;

        private final Object LOCK = new Object();

        public boolean isDone() {
            synchronized(LOCK) {
                do {
                    System.out.println("Queued = " + queued + "; queue size = " + queue.size());
                    if (!queue.isEmpty())
                        return false;
                    if (queued == 0)
                        return true;
                    try {
                        LOCK.wait();
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                } while (true);
            }
        }

        public boolean hasNext() {
            return !queue.isEmpty();
        }

        public Object get() throws InterruptedException {
            synchronized(LOCK) {
                Object res = queue.take();
                queued++;
                LOCK.notifyAll();
                return res;
            }
        }

        public void put(Object o) {
            if (queue.offer(o))
                synchronized(LOCK) {
                    LOCK.notifyAll();
                }
        }

        public void jobDone() {
            synchronized(LOCK) {
                queued--;
                LOCK.notifyAll();
            }
        }
    }

    private static class Consumer implements Runnable {

        private long startTime;

        private Producer producer;

        public Consumer(Producer producer) {
            this.producer = producer;
            startTime = System.currentTimeMillis();
        }

        public void run() {
            while (!producer.isDone()) {
                try {
                    System.out.println(producer.get());
                    System.out.println(Thread.currentThread().getId());
                    Thread.yield();
                } catch (InterruptedException e) {

                }
                if (System.currentTimeMillis() <= startTime + 15000) {
                    //System.out.println("Putting object");
                    producer.put(new Object());
                    try {
                        Thread.sleep(50);
                    } catch (InterruptedException e) {
                        
                    }
                    //System.out.println("Putting object");
                    producer.put(new Object());
                }
                producer.jobDone();
            }
            try {
                Thread.sleep(100);
            } catch (InterruptedException e) {
                
            }
            System.out.println("Finished");
        }
    }
}


Как обычно я не знаю как это работает (в смысле, синхронизация на додумках и соплях), но работает. Пугает меня метод isDone в классе Producer. Работает он верно, но ответ выдавать не сразу, а только когда поймет, завершена ли работа или нет. Правильно ли так делать? Если думать с точки зрения удобства, считаю, удобно, программисту, разработчику Consumer не надо думать и гадать, как сделать правильное ожидание разрешения ситуации.

И я что-то запамятовал, почему этот код работает?
Код

while (p.hasNext()) {
    Object o = p.get();
}

hasNext  и get синхронизированы. Но если подумать hasNext могут запросить 2 потока, сначала одному будет сказано, что есть следующий, потом второму будет сказано, что есть следующий, а затем они вдвоем помчатся ловить p.get(), один из них может навечно остаться ждать данных в p.get()

Добавлено через 12 минут и 54 секунды
Может правильней такая конструкция:

Код

while (true) {
    synchronized(producer) {
        if (producer.isDone()) break;
        System.out.println(producer.get());
    }
}


Тоже работает.

Это сообщение отредактировал(а) Platon - 25.9.2008, 10:57
PM MAIL ICQ   Вверх
Platon
Дата 25.9.2008, 18:07 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Эксперт
***


Профиль
Группа: Завсегдатай
Сообщений: 1801
Регистрация: 25.4.2006

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



В общем, уважаемые.

Решение написал сам себе. Без проблем принимаю критику.

Код

package ru.vingrad.platon.bfs.multithread;

public class BFSConsumer<T> implements Runnable {

    private final BFSQueue<T> producer;

    private final ConsumerCallback<T> callback;

    private static final boolean DEBUG = true;

    public BFSConsumer(BFSQueue<T> producer, ConsumerCallback<T> callback) {
        this.producer = producer;
        this.callback = callback;
    }

    public void run() {
        while (true) {
            T link;
            try {
                synchronized(producer) {
                    if (producer.isDone()) break;
                    link = producer.get();
                }
                callback.write(link);
            } catch (InterruptedException e) {
                break;
            }
            producer.jobDone(link);
        }
        if (DEBUG)
            System.out.println("Finished");
    }

    public static interface ConsumerCallback<T> {
        void write(T item);
    }
}


Код

package ru.vingrad.platon.bfs.multithread;

import java.util.concurrent.ArrayBlockingQueue;

public class BFSQueue<T> {

    private static final boolean DEBUG = true;

    protected QueueSource<T> queue;

    protected int queued;

    protected final Object LOCK = new Object();

    public BFSQueue(QueueSource<T> queue) {
        this.queue = queue;
    }


    public BFSQueue() {
        this(new ArrayBlockingQueueSource<T>());
    }

    public boolean isDone() {
        synchronized(LOCK) {
            do {
                if (DEBUG)
                    System.out.println("Queued = " + queued + "; queue size = " + queue.size());
                if (!queue.isEmpty())
                    return false;
                if (queued == 0)
                    return true;
                try {
                    LOCK.wait();
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            } while (true);
        }
    }

    public T get() throws InterruptedException {
        synchronized(LOCK) {
            T res = queue.get();
            queued++;
            LOCK.notifyAll();
            return res;
        }
    }

    public void put(T o) throws InterruptedException {
        synchronized(LOCK) {
            queue.put(o);
            LOCK.notifyAll();
        }
    }

    public void jobDone(T item) {
        synchronized(LOCK) {
            queued--;
            queue.done(item);
            LOCK.notifyAll();
        }
    }

    public static interface QueueSource<T> {
        T get() throws InterruptedException;
        void put(T item) throws InterruptedException;
        boolean isEmpty();
        int size();
        void done(T item);
    }

    private static class ArrayBlockingQueueSource<T> implements QueueSource<T> {

        private ArrayBlockingQueue<T> queue = new ArrayBlockingQueue<T>(4096);

        public T get() throws InterruptedException {
            return queue.take();
        }

        public void put(T item) throws InterruptedException {
            queue.put(item);
        }

        public boolean isEmpty() {
            return queue.isEmpty();
        }

        public int size() {
            return queue.size();
        }


        public void done(T item) {

        }
    }
}

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

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

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


 




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


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

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