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

Поиск:

Ответ в темуСоздание новой темы Создание опроса
> multithreading NIO Server 
V
    Опции темы
Tony
  Дата 26.11.2006, 01:12 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Эксперт
***


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

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



Всем привет.
Пишу сервер на НИО. Хочу на каждий запрос создавать новый Thread, 4то бы не ставить клиентов в очередь. Так вот, почемуто создаётся эхо соединения. Тоесть если клиент выслал 1 строчку то сервер по4емуто вызовет READ около 10(может варьироваться) раз. Причём 9 Тредов будер с буффером 0 ,а вот последний  раз нормально прочитает буффер будет равен количеству висланых байт. Если убрать из new ServerCommands(sc,keys,key).start(); метод start() то всё работает нормально. Вот код сервера.
Код

try{
            ServerSocketChannel ssc = ServerSocketChannel.open();
            ssc.configureBlocking(false);
            ServerSocket ss = ssc.socket();
            ss.bind(new InetSocketAddress(PORT));
            
            Selector selector = Selector.open();
            ssc.register(selector, SelectionKey.OP_ACCEPT );
            System.out.printf("%s %s","Server  startup on port: ",PORT);
            while(true){
                if(selector.select()>0){
                    Set <SelectionKey> keys = selector.selectedKeys();
                    Iterator<SelectionKey> it = keys.iterator();
                    
                    while(it.hasNext()){                
                        SelectionKey key = it.next();                
                            it.remove();
                            
                            SocketChannel sc = null;
                            if (key.readyOps()==SelectionKey.OP_ACCEPT){
                                Socket s = ss.accept();
                                System.out.println("ACCEPT");
                                System.out.println(s.getRemoteSocketAddress());
                                sc = s.getChannel();
                                sc.configureBlocking(false);                    
                                sc.register(selector,SelectionKey.OP_READ);
                            }else{
                                try{
                                if(key.readyOps()==SelectionKey.OP_READ){
                                    System.out.println("READ - ");
                                    
                                    sc = (SocketChannel)key.channel();
                                    new ServerCommands(sc,keys,key).start();
                                    
                                    
                                }
                                }catch(Exception e){
                                    e.printStackTrace();
                                    key.cancel();
                                    sc.close();
                                }
                        }
                        
                    }
                    keys.clear();
                }
                
            }
        }catch(Exception e){
            e.printStackTrace();
        }



--------------------
user posted image
user posted image
PM MAIL Skype   Вверх
KOp4iK
Дата 27.11.2006, 13:57 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Шустрый
*


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

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



на сколько я могу предположить 

Код

key.readyOps()==SelectionKey.OP_READ

true до тех пор пока ты из него не прочитаешь все пришедшие в него данные... 
Вот и получается что до тех пор пока ты прочитал все данные ты успеваешь создать 9-10 тредов... 
Решение как мне кажется: читать в одном треде и только потом передавать данные в тред на выполнение...
PM MAIL   Вверх
Tony
Дата 27.11.2006, 21:04 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Эксперт
***


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

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



Ведь в коде есть it.remove(); . Поэтому повторной итэрации быть не может. Распиши поподробней пожалуйста.

Это сообщение отредактировал(а) Tony - 27.11.2006, 21:05


--------------------
user posted image
user posted image
PM MAIL Skype   Вверх
Tony
Дата 29.11.2006, 13:35 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Эксперт
***


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

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



 smile 


--------------------
user posted image
user posted image
PM MAIL Skype   Вверх
iLoveJava
Дата 29.7.2007, 22:08 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Новичок



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

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



У меня тоже вопрос по этой теме.

Какая вообще должна быть структура сервера?  smile 

Код

try{
            ServerSocketChannel ssc = ServerSocketChannel.open();
            ssc.configureBlocking(false);
            ServerSocket ss = ssc.socket();
            ss.bind(new InetSocketAddress(PORT));
            
            Selector selector = Selector.open();
            ssc.register(selector, SelectionKey.OP_ACCEPT );

            while(true){
                if(selector.select()>0){
                    Set <SelectionKey> keys = selector.selectedKeys();
                    Iterator<SelectionKey> it = keys.iterator();
                    
                    for(;it.hasNext();){                
                        SelectionKey key = it.next();                
                            it.remove();
                            
                            SocketChannel sc = null;
                            if ( key.isAcceptable() ){
                                   ...
                            }else{
                                   if( key.isReadable() ){
                                         ....
                                   }
                                         ....
                            }
                        
                    }
                    keys.clear();
                }
                
            }
        }catch(Exception e){
            e.printStackTrace();
        }


Тоесть 1 главный поток который просыпается только при возникновении соответствующих событий и в нем происходит обработка сообщений иначе будет как говорил, KOp4iK
Цитата

на сколько я могу предположить 
код Java
1:
  key.readyOps()==SelectionKey.OP_READ




true до тех пор пока ты из него не прочитаешь все пришедшие в него данные... 
Вот и получается что до тех пор пока ты прочитал все данные ты успеваешь создать 9-10 тредов... 
Решение как мне кажется: читать в одном треде и только потом передавать данные в тред на выполнение...

а если много клиентов? мне кажется один поток просто не успеет...   smile smile  smile  

У меня сейчас что то типа
Код

public class Server implements IServer, Runnable {
    private class ThreadReadListener implements Runnable {
        private Selector selector = null;
        private IServerListener listener = null;

        public ThreadReadListener() {
            try {
                this.listener = ServerListenerFactory.newServerListener();
                selector = Selector.open();
            } catch (IOException e) {
                // TODO Auto-generated catch block
                e.printStackTrace();
            } catch (Exception e1) {
                // TODO Auto-generated catch block
                e1.printStackTrace();
            }
        }
        
        public void run() {
            Iterator iterator = null; 
            SelectionKey skey = null;
            SocketChannel channel = null;
            for(;;) {
                try {
                    selector.select(SELECT_SLEEP_TIME);
                    iterator = selector.selectedKeys().iterator();

                    for(;iterator.hasNext();) {
                     skey = (SelectionKey) iterator.next();
                     iterator.remove();
                     if(
                         ((skey.readyOps() & SelectionKey.OP_READ) != 0) &&
                         ((skey.readyOps() & SelectionKey.OP_WRITE) != 0)
                     ) {
                         channel = (SocketChannel)skey.channel();
                         channel.configureBlocking(false);
                         listener.onRecive( channel );
                     }
                     skey.isReadable()
                     if( !skey.isValid() )
                         listener.onDisconnect(channel);
                    }
                    selector.selectedKeys().clear();
                } catch (IOException e) {
                    // TODO Auto-generated catch block
                    e.printStackTrace();
                }
            }
        }

        public void addChanel(SocketChannel chanel) throws IOException {
            chanel.configureBlocking(false);
            chanel.register(selector, SelectionKey.OP_READ | SelectionKey.OP_WRITE);
        }    
    }
    
    private static final int THREAD_COUNT = 20;
    private static final int SELECT_SLEEP_TIME = 500;

    private static Integer listenerNum = 0;
    private boolean run = false;
    private ServerSocketChannel serverChannel = null;
    private IServerListener acceptListener = null;
    private Thread threadAcceptListener = null;
    private ThreadReadListener readers[] = new ThreadReadListener[THREAD_COUNT];
    private Thread readerThreads[] = new Thread[THREAD_COUNT];
    
    public boolean isRun() {
        // TODO Auto-generated method stub
        return run;
    }

    public void start(int port) throws IOException, Exception {
        serverChannel = ServerSocketChannel.open();
        serverChannel.socket().bind( new InetSocketAddress(port) );
        
        for(int i = 0; i < THREAD_COUNT; i++) {
            readers[i] = new ThreadReadListener();
            readerThreads[i] = new Thread( readers[i] );
            readerThreads[i].start();
        }
        
        acceptListener = ServerListenerFactory.newServerListener();
        
        threadAcceptListener = new Thread(this);
        threadAcceptListener.start();
    }

    public void stop() throws IOException {
        threadAcceptListener.stop();
        for(int i = 0; i < THREAD_COUNT; i++)
            readerThreads[i].stop();
        serverChannel.close();
        run = false;
    }

    public void run() {
        run = true;
        
        SocketChannel chanel = null;
        for(;;) {
            try {
                chanel = serverChannel.accept();
                acceptListener.onConnect( chanel );
                synchronized (listenerNum) {
                    readers[listenerNum].addChanel(chanel);
                    listenerNum += 1;
                    if( listenerNum >= readers.length )
                        listenerNum = 0;
                }
            } catch (IOException e) {
                // TODO Auto-generated catch block
                e.printStackTrace();
            }
        }
    }
}


Тоесть 1 главный поток который просыпается только при новом подключении, и добавляет его в selector 1 из N потоков, которые уже ждут пока в канале не появятся данные, а потом считывают их и выполняют обработку, по моему и то лучше... НО КАК СДЕЛАТЬ ПРАВЕЛЬНО?!!!  smile 

И еще один вопрос...
Мне требуется не только читать данные но и писать по тому же запросу? как это сделать?  smile 
у меня сейчас так
Код

selector.select();
iterator = selector.selectedKeys().iterator();

for(;iterator.hasNext();) {
     skey = (SelectionKey) iterator.next();
     iterator.remove();
     if( ((skey.readyOps() & SelectionKey.OP_READ) != 0) &&
    ((skey.readyOps() & SelectionKey.OP_WRITE) != 0)
    ) {
          сhannel = (SocketChannel)skey.channel();
     channel.configureBlocking(false);
     listener.onRecive( channel );
     }

     // как мне узнать что сокет уже закрытый
     if( !skey.isValid() )
          listener.onDisconnect(channel);
}

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


Новичок



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

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



Насклько я понял читая топик

Nio, select, длиные запросы, сохранение канала, Как реализовать оптимально? 

структура сервера должна быть такой(я там тоже это написал может хоть кто-то ответит)
должен быть цикл с select в, котором для каналов будут считываться и записываться данные,
при чем у каждого канала должен быть буфер на чтение и на запись, в этом цикле буфер на чтение будет заполняться, а потом когда накопится все сообщение, то в каком-то рабочем потоке(их наверно должен быть пул) этот буфер обрабатывается и заполняется буфер на запись, который в основном цикле записуется. так чтоли?
PM MAIL   Вверх
COVD
Дата 30.7.2007, 15:16 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Эксперт
***


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

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



В сети много примеров нио сервера. Мне понравился http://rox-xmlrpc.sourceforge.net/niotut/index.html .

Цитата

должен быть цикл с select в, котором для каналов будут считываться и записываться данные, при чем у каждого канала должен быть буфер на чтение и на запись,


да


 
Цитата

в этом цикле буфер на чтение будет заполняться, а потом когда накопится все сообщение, то в каком-то рабочем потоке(их наверно должен быть пул) этот буфер обрабатывается и заполняется буфер на запись, который в основном цикле записуется. так чтоли?


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

В асинхронной же схеме клиент может прислать подряд несколько запросов. А сервер будет отвечать по мере готовности ответов, причем ответы могут идти в  другом порядке (потому что запросы могут быть разными по тяжести  - на одни ответ быстро формируется, на другой медленно). Кроме того, в асинхронной схеме сервер может что-то послать клиенту без запроса, т.е. сделать ему push.  
Поэтому в асинхронной схеме поток, обслуживающий select , всегда выполняет две операции - чтение и запись:
1. если канал готов к чтению (клиент что-то прислал, хотя бы один байт) , то читает из канала доступные байты в клиентский буфер на чтение.
2. если канал готов к записи и клиентский буфер на запись не пустой (т.е. рабочие потоки, обрабатывающие клиентские запросы, приготовили ответы), то пишет в канал из буфера на запись.

Это сообщение отредактировал(а) COVD - 30.7.2007, 15:17
PM MAIL   Вверх
iLoveJava
Дата 31.7.2007, 13:09 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Новичок



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

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



А разве в такой схеме так не выходит?

1) Приходит 1 запрос, он считывается в основном цикле, потом отдается на выполнение 1 с рабочих потоков, которому на это требуется 5минут;
2) В это время приходит 2 запрос, он считывается в основном цикле, потом отдается на выполнение 2 с рабочих потоков, которому на это требуется 2минут;
3) 2 поток завершает обработку запроса и пишет результат в буфер для отправки данных, с которого потом в основном цикле будут писаться данные в канал;
4) 1 поток завершает обработку запроса и пишет результат в буфер для отправки данных, с которого потом в основном цикле будут писаться данные в канал;

параллельно всему этому в основном потоке проверяется есть ли данные на запись если есть то они пишутся.

В результате ответ на запрос 2 приходит раньше чем на запрос 1


У меня тут еще один вопрос появился  smile 
Если я регистрирую канал на чтение и запись, то выходит, что select будет завершаться очень быстро так как канал, практически всегда будет готов на запись, то есть если сервер пишет в канал только результаты запросов, то пока запрос не пришел нет смысла проверять канал на запись в этом есть смысл только если уже пришел запрос(если конечно сервер сам может что-то записать в любой момент клиенту то тогда смысл есть, но все же если эта операция делается очень редко, а даже если часто всеравно не рационально), то есть выйдет что у меня все время будет перебор каналов на запись, а мне туда и писать нечего, как это обойти?
Можно канал зарегистрировать сначала на чтение, а когда отправлять на обработку ответа зарегистрировать на запись или после обработки зарегистрировать на запись, а после отсылки всех данных зарегистрировать на чтение(а в таком случае не может быть проблем с асинхронной схемой?)?
 smile   smile  smile 
PM MAIL   Вверх
COVD
Дата 31.7.2007, 16:02 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Эксперт
***


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

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



Цитата

В результате ответ на запрос 2 приходит раньше чем на запрос 1


Да, вы правильно описали асинхронную модель. А реализуется ли это в вашем коде - не знаю.


Цитата

Если я регистрирую канал на чтение и запись, то выходит, что select будет завершаться очень быстро так как канал, практически всегда будет готов на запись, то есть если сервер пишет в канал только результаты запросов, то пока запрос не пришел нет смысла проверять канал на запись в этом есть смысл только если уже пришел запрос(если конечно сервер сам может что-то записать в любой момент клиенту то тогда смысл есть, но все же если эта операция делается очень редко, а даже если часто всеравно не рационально), то есть выйдет что у меня все время будет перебор каналов на запись, а мне туда и писать нечего, как это обойти?


Совершенно верно. Про это как раз и написано в ссылке, которую я дал выше. 

Цитата

Set OP_WRITE only when you have data ready

A common mistake is to enable OP_WRITE on a selection key and leave it set. This results in the selecting thread spinning because 99% of the time a socket channel is ready for writing. In fact the only times it's not going to be ready for writing is during connection establishment or if the local OS socket buffer is full. The correct way to do this is to enable OP_WRITE only when you have data ready to be written on that socket channel. And don't forget to do it from within the selecting thread. 


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

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

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


 




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


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

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