Версия для печати темы
Нажмите сюда для просмотра этой темы в оригинальном формате
Форум программистов > Java: Общие вопросы > Программа с несколькими потоками и семафорами


Автор: libertas 9.4.2014, 08:55
Всем привет!

Написал программу, где есть несколько потоков - читатели и писатели.
Есть воображаемая БД. Когда читатель хочет что- нибудь прочитать, он проверяет не пишет ли писатель сейчас данные в БД. Если пишет - то читатель ожидает.
Аналогично писатель - ожидает пока все читатели прочитают.
Читателей может быть несколько одновременно, писатель - только один.

Нужно задачку решить через семафоры, которые нужно написать самому.


Код

public class MyProg {
    public static void main(String[] args) {
        Database myDB = new Database();
        Reader reader1 = new Reader(1, myDB);
        Reader reader2 = new Reader(2, myDB);
        Reader reader3 = new Reader(3, myDB);
        Reader reader4 = new Reader(4, myDB);
        Reader reader5 = new Reader(5, myDB);
        Writer writer1 = new Writer(1, myDB);
        Writer writer2 = new Writer(2, myDB);
        Writer writer3 = new Writer(3, myDB);
        Writer writer4 = new Writer(4, myDB);
        reader1.start();
        reader2.start();
        reader3.start();
        reader4.start();
        reader5.start();
        writer1.start();
        writer2.start();
        writer3.start();
        writer4.start();
    }
}
 
 
class Reader extends Thread {
    public Reader(int r, Database db) {
        readerNum = r;
        server = db;
    }
 
    public void run() {
        int c;
        for (int i = 1; i <= 3; i++) 
        {
            napping();
            System.err.println("reader " + readerNum + " wants to read.");
 
            c = server.startRead();
            System.err.println("reader " + readerNum + " is reading. Reader Count = " + c);
            Database.something();
 
            c = server.endRead();
        }
    }
 
    public void napping() {
        try {
            sleep((int) (Math.random() * 5000));
        } catch (InterruptedException e) {
        }
    }
 
    private Database server;
    private int readerNum;
}
 
 
class Writer extends Thread {
    public Writer(int w, Database db) {
        writerNum = w;
        server = db;
    }
 
    public void run() {
        for (int i = 1; i <= 3; i++) 
        {
            napping();
            System.err.println("writer " + writerNum + " wants to write.");
            server.startWrite();
 
            System.err.println("-------------- WRITER " + writerNum + " IS WRITING.");
            Database.something();
 
            System.err.println("               writer " + writerNum + " is done writing.");
            server.endWrite();
 
        }
    }
 
    public void napping() {
        try {
            sleep((int) (Math.random() * 5000));
        } catch (InterruptedException e) {
        }
    }
 
    private Database server;
    private int writerNum;
}
 
 
class Database {
    public Database() {
        readerCount = 0;
        dbReading = false;
        dbWriting = false;
        Sem = new Semaphore(1);
    }
 
    public static void something() {
        try {
            Thread.sleep((int) (Math.random() * 5000));
        } catch (InterruptedException e) {
        }
    }
 
     public int startRead() {
        Sem.P();
        while (dbWriting) {
            try {
                wait();
            } catch (InterruptedException e) {
            }
        }
        ++readerCount;
        if (readerCount == 1)  // if I am the first reader tell all others
            dbReading = true;    // that the DB is being read
        Sem.V();
        return readerCount;
    }
 
    public int endRead() {
        Sem.P();
        --readerCount;
        if (readerCount == 0)  // if I am the last reader tell all others
            dbReading = false;  // that the database is no longer being read
        System.err.println("Reader is done reading. Count = " + readerCount);
        Sem.V();
        return readerCount;
    }
 
    public void startWrite() {
        Sem.P();
        while (dbReading || dbWriting) {
            try {
                wait();
            } catch (InterruptedException e) {
            }
        }
        dbWriting = true; // the DB is being written
        Sem.V();
    }
 
    public void endWrite() {
        Sem.P();
        dbWriting = false;
        Sem.V();
    }
 
    private int readerCount;    // the number of active readers
    private boolean dbReading;  // flags to indicate whether the DB
    private boolean dbWriting;  // is being read or written
    private Semaphore Sem;
}
 
class Semaphore {
    private int value;
 
    Semaphore(int v) {
        value = v;
    }
 
    public void P() {
        while (value <= 0) ;
        value--;
    }
 
    public void V() {
        ++value;
    }
}


Я написал семафор - Semaphore, который в зависимости от значения value проверяет - читать / писать.
Я понимаю, что эти семафоры не атомарны, но задачку надо решить именно через семафоры, написанные вручную.

В этом коде проиходит исключение:

Exception in thread "Thread-3" java.lang.IllegalMonitorStateException
at java.lang.Object.wait(Native Method)
at java.lang.Object.wait(Object.java:503)
at Database.startRead(MyProg.java:113)
at Reader.run(MyProg.java:39)

В чем может быть дело?

Спасибо

Автор: LSD 9.4.2014, 19:03
Ошибка говорит о том, что метод wait() должен вызываться потоком которых захватил монитор на объект (у которого вызыватся wait()). Но это все не важно, тут есть проблемы серьезней:
1. Непонятно причем тут семафоры, когда это типичный read-write lock.
2. Реализовать семафор самостоятельно без "нативной" поддержки нельзя. Или надо использовать synchronized или готовую реализацию из java.util.concurrent.

Автор: Mirkes 9.4.2014, 22:33
Цитата(LSD @  9.4.2014,  19:03 Найти цитируемый пост)
Реализовать семафор самостоятельно без "нативной" поддержки нельзя. 

Пожалуй рискну не согласиться. Реализовать семафор все-таки можно, но семафор нужно опрашивать на предмет разрешенности действия. А вот опроса семафора я в коде не увидел. По идее должно быть что-то вроде: Читатель семафору "Можно читать?", Семафор читателю "Можно", тогда стартует чтение. Правда при многопоточности нет гарантии, что между ответом семафора и собственно началом чтения не влезет кто-то еще. Для предотвращения такого конфликта, семафор в случае положительного ответа сначала блокирует доступ другим, а уже затем выдает ответ. Получается кривовато, но работоспособно, если читатель не забудет читать после запроса и т.д.
Но в предположении приличного поведения читателей и писателей семафор сделать можно.

В реальное, не учебное, приложение я бы такой семафор не включал  smile 

Автор: libertas 10.4.2014, 07:49
Цитата(LSD @  9.4.2014,  19:03 Найти цитируемый пост)
Ошибка говорит о том, что метод wait() должен вызываться потоком которых захватил монитор на объект (у которого вызыватся wait()). Но это все не важно, тут есть проблемы серьезней:
1. Непонятно причем тут семафоры, когда это типичный read-write lock.
2. Реализовать семафор самостоятельно без "нативной" поддержки нельзя. Или надо использовать synchronized или готовую реализацию из java.util.concurrent. 


Спасибо. 
Да согласен, именно с этим и была проблема. Плюс проблема в растановке семафоров.

Вот переделанный рабочий код:

Код

public class MyProg {
    public static void main(String[] args) {
        Database myDB = new Database();
        Reader reader1 = new Reader(1, myDB);
        Reader reader2 = new Reader(2, myDB);
        Reader reader3 = new Reader(3, myDB);
        Reader reader4 = new Reader(4, myDB);
        Reader reader5 = new Reader(5, myDB);
        Writer writer1 = new Writer(1, myDB);
        Writer writer2 = new Writer(2, myDB);
        Writer writer3 = new Writer(3, myDB);
        Writer writer4 = new Writer(4, myDB);
        reader1.start();
        reader2.start();
        reader3.start();
        reader4.start();
        reader5.start();
        writer1.start();
        writer2.start();
        writer3.start();
        writer4.start();
    }
}


class Reader extends Thread {
    public Reader(int r, Database db) {
        readerNum = r;
        server = db;
    }

    public void run() {
        int c;
        for (int i = 1; i <= 3; i++) // It should be while (true)
        {
            napping();
            System.err.println("reader " + readerNum + " wants to read.");

            c = server.startRead();
            System.err.println("reader " + readerNum + " is reading. Reader Count = " + c);
            Database.something();

            c = server.endRead();
        }
    }

    public void napping() {
        try {
            sleep((int) (Math.random() * 5000));
        } catch (InterruptedException e) {
        }
    }

    private Database server;
    private int readerNum;
}


class Writer extends Thread {
    public Writer(int w, Database db) {
        writerNum = w;
        server = db;
    }

    public void run() {
        for (int i = 1; i <= 3; i++) // It should be while (true)
        {
            napping();
            System.err.println("writer " + writerNum + " wants to write.");
            server.startWrite();

            System.err.println("-------------- WRITER " + writerNum + " IS WRITING.");
            Database.something();

            System.err.println("               writer " + writerNum + " is done writing.");
            server.endWrite();
        }
    }

    public void napping() {
        try {
            sleep((int) (Math.random() * 5000));
        } catch (InterruptedException e) {
        }
    }

    private Database server;
    private int writerNum;
}


class Database {
    public Database() {
        readerCount = 0;
        semWriter = new Semaphore(1);
        sem = new Semaphore(1);
    }

    public static void something() {
        try {
            Thread.sleep((int) (Math.random() * 5000));
        } catch (InterruptedException e) {
        }
    }

    public int startRead() {
        sem.P();
        if (readerCount == 0){
            semWriter.P();
        }
        readerCount++;
        sem.V();
        return readerCount;
    }

    public int endRead() {
        sem.P();
        --readerCount;
        System.err.println("Reader is done reading. Count = " + readerCount);
        if (readerCount == 0){
            semWriter.V();            // if I am the last reader tell all others
        }
        sem.V();
        return readerCount;
    }

    public void startWrite() {
        semWriter.P();
    }

    public void endWrite() {
        semWriter.V();
    }

    private int readerCount;    // the number of active readers
    private Semaphore semWriter, sem;
}

class Semaphore
{
    public Semaphore(int v)
    { value = v; }

    public synchronized void P()
    { value--;
        if (value < 0)
            try
            {wait();}
            catch (InterruptedException e) {}
    }

    public synchronized void V()
    { ++value;
        if (value<=0)
            notify();
    }

    private int value;
}



Добавлено через 2 минуты и 5 секунд
Цитата(Mirkes @  9.4.2014,  22:33 Найти цитируемый пост)
 По идее должно быть что-то вроде: Читатель семафору "Можно читать?", Семафор читателю "Можно", тогда стартует чтение. Правда при многопоточности нет гарантии, что между ответом семафора и собственно началом чтения не влезет кто-то еще.


Согласен.

Цитата(Mirkes @  9.4.2014,  22:33 Найти цитируемый пост)
 
 Для предотвращения такого конфликта, семафор в случае положительного ответа сначала блокирует доступ другим, а уже затем выдает ответ. 


Извините, не очень понятно как это можно сделать?

Автор: LSD 10.4.2014, 10:49
Цитата(Mirkes @  9.4.2014,  23:33 Найти цитируемый пост)
Пожалуй рискну не согласиться. Реализовать семафор все-таки можно, но семафор нужно опрашивать на предмет разрешенности действия. А вот опроса семафора я в коде не увидел. По идее должно быть что-то вроде: Читатель семафору "Можно читать?", Семафор читателю "Можно", тогда стартует чтение. Правда при многопоточности нет гарантии, что между ответом семафора и собственно началом чтения не влезет кто-то еще. Для предотвращения такого конфликта, семафор в случае положительного ответа сначала блокирует доступ другим, а уже затем выдает ответ. Получается кривовато, но работоспособно, если читатель не забудет читать после запроса и т.д.
Но в предположении приличного поведения читателей и писателей семафор сделать можно.

1. Не решается вопрос как "усыпить" поток на время ожидания захвата семафора. Можно конечно busy spin сделать, но это не очень хорошее решение в общем случае.
2. Без блокировок или CAS нельзя произвести атомарный инкремент/декремент (хватит ли одного CAS для реализации семафора я не уверен).


Цитата(libertas @  10.4.2014,  08:49 Найти цитируемый пост)
Вот переделанный рабочий код:

Твой семафор работает неправильно. В случае если счетчик на нуле, и 2 потока вызовут acquire(), счетчик станет -2. Потом один поток вызовет release() счетчик станет -1 и один из потоков пробудится, что неправильно.

Автор: libertas 10.4.2014, 11:06
Цитата(LSD @  10.4.2014,  10:49 Найти цитируемый пост)
Твой семафор работает неправильно. В случае если счетчик на нуле, и 2 потока вызовут acquire(), счетчик станет -2. Потом один поток вызовет release() счетчик станет -1 и один из потоко пробутидся, что неправильно. 


Извините, можно пояснить.

Где/когда 2 моих  потока будут вызывать acquire()? И где/когда один мой поток будет вызывать release()?

Автор: LSD 10.4.2014, 14:03
Цитата(libertas @  10.4.2014,  12:06 Найти цитируемый пост)
Где/когда 2 моих  потока будут вызывать acquire()? И где/когда один мой поток будет вызывать release()?

У тебя там куча ошибок:
1. Операция ++ не атомарна. В многопоточной среде может давать неправильный результат. Плюс у тебя переменные не volatille потоки увидят изменения сделанные друг другом в непредсказуемый момент.
2. Операция
Код

if (<check condition>)
    doOperation();

не атомарна. Возможна ситуация: 
- проверили условие оно истинно
- ОС усыпила поток
- другой поток поменял что-то и условие уже не выполняется
- ОС разбудила первый поток и он пошел выполнять doOperation(), хотя условие уже не выполняется
3. Логические ошибки в Database у тебя в начале метода startRead() вызывается sem.P(), который поддерживает только один поток. Т.е. второй читатель будет ждать, хотя читатели должны работать паралельно.
4. Реализация самого Semaphore:
- поток 1 заходит в startRead() и вызывает sem.P() и усыпляется ОС где-то в середине метода, значение value = 0
- поток 2 заходит в startRead() и вызывает sem.P() и засыпает, значение value = -1
- поток 3 заходит в startRead() и вызывает sem.P() и засыпает, значение value = -2
- поток 1 выходит из метода и вызывает sem.V(), просыпается один из потоков 2 или 3, значение value = -1
В принципе котракт не нарушен - только один поток выполняется, но странно.

Автор: libertas 10.4.2014, 14:32
Цитата(LSD @  10.4.2014,  14:03 Найти цитируемый пост)
У тебя там куча ошибок:
1. Операция ++ не атомарна. В многопоточной среде может давать неправильный результат. Плюс у тебя переменные не volatille потоки увидят изменения сделанные друг другом в непредсказуемый момент.


Если Вы имеете в виду :
readerCount++;

то она защищена семафором. Если поток заснул, то другой все равно не сможет её изменить.

А если про эту:

++value;

то она находится в методе synchronized - следовательно может меняться только одним потоком. Так работает synchronized  .


Цитата

2:
if (<check condition>)
    doOperation();

не атомарна. Возможна ситуация: 
- проверили условие оно истинно
- ОС усыпила поток
- другой поток поменял что-то и условие уже не выполняется
- ОС разбудила первый поток и он пошел выполнять doOperation(), хотя условие уже не выполняется


Это невозможно, потому что
if (<check condition>)
    doOperation();
защищено synchronized  семафором.

Цитата

4. Реализация самого Semaphore:
- поток 1 заходит в startRead() и вызывает sem.P() и усыпляется ОС где-то в середине метода, значение value = 0
- поток 2 заходит в startRead() и вызывает sem.P() и засыпает, значение value = -1
- поток 3 заходит в startRead() и вызывает sem.P() и засыпает, значение value = -2
- поток 1 выходит из метода и вызывает sem.V(), просыпается один из потоков 2 или 3, значение value = -1


Тоже не может быть, потому что sem.P() synchronized .

Есть там, конечно не точность в коде. Сейчас только увидел:
Код

sem.P();
        if (readerCount == 0){
            semWriter.P();
        }
        readerCount++;
        sem.V();


В этом участке, если читатель пройдет семафор и уснет, то другие читатель не смогут заходить в этот блок, пока не проснется этот читатель и не откроет семафор.


Автор: LSD 10.4.2014, 16:13
Замечание насчет vollatile все равно остается в силе.

Все замечания про защиту монитром верны до тех пор пока у него count = 1, ты его используешь как обычный мьютекс. Проще уж сразу было делать synchronized и не мучаться. 


Цитата(libertas @  10.4.2014,  15:32 Найти цитируемый пост)
Тоже не может быть, потому что sem.P() synchronized .

Речь не про sem.P(), а про вызывающий метод.

Автор: Mirkes 10.4.2014, 20:14
Цитата(LSD @  10.4.2014,  10:49 Найти цитируемый пост)
1. Не решается вопрос как "усыпить" поток на время ожидания захвата семафора. Можно конечно busy spin сделать, но это не очень хорошее решение в общем случае.
2. Без блокировок или CAS нельзя произвести атомарный инкремент/декремент (хватит ли одного CAS для реализации семафора я не уверен).


Хорошее решение должно опираться на синхронизацию или конкуренцию (со второй не знаком пока). Я же написал, что в приличную программу я бы предложенный мной семафор не включил. 
"Усыпление" потока осуществляется тупо - он просится, пока не получит разрешения. Это очень расточительно по ресурсам, но осуществимо. Довольно давно под DOS мне приходилось реализовывать многозадачность. Это куча проблем, жестоко карается любая ошибка, но в принципе осуществимо. Единственное, что мне не понятно, зачем давать такое как учебную задачу.

Автор: libertas 11.4.2014, 07:48
Цитата(LSD @  10.4.2014,  16:13 Найти цитируемый пост)
Замечание насчет vollatile все равно остается в силе.

Все замечания про защиту монитром верны до тех пор пока у него count = 1, ты его используешь как обычный мьютекс. Проще уж сразу было делать synchronized и не мучаться. 


да, про volatile нужно почитать. Ну да тут мьютекс получается. Согласен, что проще делать быть через synchronized, просто задание было сделать именно через семафоры.

Добавлено через 37 секунд
Цитата(Mirkes @  10.4.2014,  20:14 Найти цитируемый пост)
"Усыпление" потока осуществляется тупо - он просится, пока не получит разрешения.



Ну да, ресурсы цпу жрет - крутится в цикле.

Автор: LSD 14.4.2014, 12:48
Цитата(Mirkes @  10.4.2014,  21:14 Найти цитируемый пост)
"Усыпление" потока осуществляется тупо - он просится, пока не получит разрешения. Это очень расточительно по ресурсам, но осуществимо.

Это и есть busy spin, более того он иногда оправдае, например java.util.concurrent.atomic его используют. Основная проблема во втором пункте:
Цитата(LSD @  10.4.2014,  11:49 Найти цитируемый пост)
2. Без блокировок или CAS нельзя произвести атомарный инкремент/декремент (хватит ли одного CAS для реализации семафора я не уверен).

Я немного посмотрел - CAS хватит для реализации семафора с busy spin.

Powered by Invision Power Board (http://www.invisionboard.com)
© Invision Power Services (http://www.invisionpower.com)