Страницы

Поиск по вопросам

Показаны сообщения с ярлыком сокеты. Показать все сообщения
Показаны сообщения с ярлыком сокеты. Показать все сообщения

среда, 5 февраля 2020 г.

Что происходит когда сокет записывает данные в то время, когда с другой стороны читаются предыдущие данные?

#cpp #qt5 #tcp #сокеты


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


Ответы

Ответ 1



Если сокет блокирующий и на чтение запрошено больше данных, чем Вы отправляете - будет ждать "второй итерации". Если сокет неблокирующий - вычитает, сколько сможет, и затем Вам снова придется опрашивать сокет. В связи с комментарием @VTT хочу дать немного более развернутый ответ. Если вызывается одна из функций read/recv/recvfrom для блокируемого сокета и при этом в буфере нет никаких данных - сокет переходит в спящее состояние до тех пор, пока не придут какие-либо данные. Достаточно одного байта. Тем не менее, мы можем задать флаг MSG_WAITALL - в этом случае мы будем ждать до тех пор, пока не будет доступно фиксированное (запрошенное нами) количество байт. В случае неблокируемого сокета - если не удовлетворены условия ввода - то ф-ия вернет управление установив ошибку EWOULDBLOCK. Если вызывается одна из функций write/send/sendto - данные копируются из буфера приложения в буфер отправки сокета. Для блокируемого сокета - если в буфере отправки недостаточно места, процесс переходит в состоянии оидания до тех пор, пока это самое место не освободится. Для неблокируемого сокета - если в буфере отправки недостаточно места - ф-ия вернет управление, установив ошибку EWOULDBLOCK. Комбинируйте и предполагайте наиболее вероятный результат того, чем же все-таки окончится посылка двух сообщений :)

Ответ 2



В общем случае, попытки записи и чтения данных будут не связаны. Вторая попытка записи может происходить совершенно независимо от поведения (или даже от наличия) партнера на другой стороне. Более того, когда вызов функции записи на первой итерации успешно вернется, нет оснований полагать, что только что отправленные данные на самом деле куда-то ушли. Если активирован алгоритм Нейгла, то только что посланные данные могут спокойно сидеть в буфере на стороне отправителя. Даже если подкрутить настройки (допустим TCP_NODELAY), то данные на момент второго вызова записи могут еще быть где-то в пути, или сидеть в буфере на принимающей стороне.

четверг, 11 июля 2019 г.

неблокирующий udp-сокет в java

Есть сервер, на котором должен работать демон (фоновый процесс, не важно как назвать) постоянно обрабатывающий массив информации. Периодически он должен проверять не пришли ли какие-то сообщения от сетевых клиентов и, если пришли, записывать их в базу. Из-за определённых ограничений клиенты отправляют серверу сообщения исключительно по udp протоколу. Если я правильно понял прочитанную литературу по работе с udp-сокетами, то выглядеть это будет примерно так:
-где-то вначале инициализируем сокет на определённом порту
try { socket = new DatagramSocket(port); } catch (SocketException e) { // Сообщение об ошибке, возможно остановка всего демона, если ничего не вышло }
-далее выполняем основную работу демона и в определённые моменты пытаемся получить информацию из сокета:
byte[] bytes = new byte[1024]; DatagramPacket p = new DatagramPacket(bytes, bytes.length); try { socket.receive(p); } catch(IOException e) { }
Насколько я понимаю, если ни один клиент не прислал пакета данных, то на данном этапе всё остановится и процесс будет ждать прихода данных на сокет. Правильно? А мне необходимо, чтоб работа продолжалась независимо от того, получены данные или нет. Если получены - хорошо - записали их в базу, если нет, ну и ладно - ещё масса работы. В C++ существует возможность перевода сокета в неблокирующий режим, выглядит это приблизительно так:
int nonBlocking = 1; if ( fcntl( handle, F_SETFL, O_NONBLOCK, nonBlocking ) == -1 ) { printf( "failed to set non-blocking socket
" ); }
для UNIX и MAC и
DWORD nonBlocking = 1; if ( ioctlsocket( handle, FIONBIO, &nonBlocking ) != 0 ) { printf( "failed to set non-blocking socket
" ); }
для WINDOWS где handle - дескриптор сокета. При таком режиме, при чтении из сокета, в случае, когда данные в нём отсутствуют,
int received_bytes = recvfrom( handle, (char*)packet_data, maximum_packet_size, 0, (sockaddr*)&from, &fromLength );
программа не остановится, а received_bytes просто будет равно нулю. Существует ли в Java возможность переводить udp-сокет в неблокирующий режим? Или необходимо опрос сокета вести в отдельном потоке и при получении данных реализовывать какую-то событийность? Неблокирующие сокеты, в моём случае, были бы хороши тем, что опрос сокета должен проводиться в определённые моменты обработки основных данных и если получена какая-то информация, то она должна быть сохранена с привязкой к определённым ключам из основных данных, в остальное время пришедшие udp-пакеты неважны, их можно просто сбрасывать. Бегло пробежавшись по результатам поисковых запросов не нашёл необходимого мне функционала в Java.

Почему поток для меня неудобен, попробую объяснить: всё это дело связано со сложной системой охранной сигнализации и видеонаблюдения. Грубо говоря: одни железки (железо1) постоянно пишут какие-то данные в базу. Один из процессов на сервере (процесс1) эту базу постоянно контролирует и при обнаружении в ней определённых новых данных рассылает другим железкам (железо2) сообщения. Те, в свою очередь, получив сообщения проводят определённые действия и проверки (это занимает какое-то небольшое время) и при возникновении у них каких-то событий отправляют данные на сервер по протоколу udp. Возникла задача получать эти сообщения (не критично, если некоторые будут теряться), сравнивать с данными из базы и при некоторых совпадениях - логировать. Существующее программное обеспечение не изменить - нет исходников. Так вот проблема в том, что железо2, присылая сообщения, не ставит там никаких временных меток. А ещё оно может самостоятельно присылать сообщения, не инициированные запросами от процесса1. И мне было бы удобно в моём процессе мониторить базу (одновременно с другим процессом, это уже реализовано и сложностей с этим не возникает) и пытаться считать данные из сокета только в определённые моменты, когда там (в базе) появляются конкретные данные (эти события происходят довольно редко, раз в два-три часа), а в остальное время просто периодически сбрасывать данные сокета (если железо2 туда что-то прислало), как ненужные. Иначе я могу получить кучу данных, присланных за последнее время, которые к нужным мне событиям отношения не имеют, но выяснить это можно лишь дифференцировав это всё по времени, а временных меток (повторюсь) в сообщениях железа2 нет Если реализовывать это в отдельном потоке, то для меня возникают определённые сложности с созданием событий и их вызовом. А неблокирующий сокет в этом случае - идеальное решение. Тем более, что такая функциональность на C++ точно есть, она работает и работает прекрасно. Вот и возник вопрос: возможно подобный функционал есть и у Java? Заодно, раз уж приходится всё делать именно на этом языке программирования, буду развиваться и учить для себя что-то новое.


Ответ

В Java все non-blocking IO содержится в пакете java.nio.channels. Неблокирующее UDP реализовано классом DatagramChannel
Пользоваться можно как-то так:
DatagramChannel chan = DatagramChannel.open(); chan.configureBlocking( false ); chan.bind( new InetSocketAddress( 7777 ) );
ByteBuffer buffer = ByteBuffer.allocate( 4*1024 );
while (true) { buffer.clear(); System.out.println("trying non-blocking receive..."); SocketAddress from = chan.receive(buffer); System.out.println("non-blocking receive done.");
if (from != null) { buffer.flip(); System.out.printf("<<<--- got [%x] byte from %s%n", buffer.get(), from); }
System.out.println( "sleeping..." ); TimeUnit.SECONDS.sleep( 5 ); }

пятница, 5 июля 2019 г.

Проблема с ObjectInputStream. (StreamCorruptedException: invalid type code: AC)

Вообще изначальная задача - написать игру с двумя игроками, что-то типа точек. т.е. один игрок (клиент) делает ход, отправляет данные (координаты x и y) на сервер. Сервер обрабатывает эту информацию и отправляет сообщение с какими-то данными всем подключенным клиентам. Все это надо реализовать на сокетах. За основу взяла реализацию чата, но мне удобнее передавать объекты, поэтому использую ObjectInputStream и ObjectOutputStream. Написала пример: клиент вводит координаты x y, отправляет на сервер. На сервере хранится матрица А. Получаем от клиента x,y, присваиваем A[x][y]=1, отправляем всем клиентам объект класса ServerAnswer (матрица, x,y). Клиент получает и выводит матрицу. Для начала запускаю только один клиент. Первый раз все работает, но когда я ввожу новые координаты на клиенте, они отправляются на сервер, он посылает ответ и на клиенте при попытке выполнения - sa = (ServerAnswer)ois.readObject(); вылетает ошибка апр 17, 2014 12:06:45 PM client2.SocketInputThread run SEVERE: null java.io.StreamCorruptedException: invalid type code: AC at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1377) at java.io.ObjectInputStream.readObject(ObjectInputStream.java:370) at client2.SocketInputThread.run(SocketInputThread.java:37) at java.lang.Thread.run(Thread.java:744) Подскажите. пожалуйста, что не так в коде. Код Сервера: public class Server2 {
public static void main(String[] args) { System.out.println("Program starting..."); try { ServerSocket ss = new ServerSocket(3129,0, InetAddress.getByName("localhost")); System.out.println("Server starting..."); while(true){ Socket s = ss.accept(); // ожидание новых клиентов SocketThread socketThread = new SocketThread(s); Thread t = new Thread(socketThread); t.start(); // запуск нового потока для каждого нового клиента } } catch (IOException ex) { Logger.getLogger(Server2.class.getName()).log(Level.SEVERE, null, ex); }
} } Класс SocketThread: public class SocketThread implements Runnable {
private Socket s = null;
private boolean exit = true; private int x = -1; private int y = -1; private int p = 0; private int [][]A; private ArrayList listSocket = null; private ObjectInputStream ois = null; private ObjectOutputStream oos = null;
public SocketThread(Socket s) { this.s = s; }
@Override public void run() { try { System.out.println("User connect..."); ListSocket.addSocketToList(s); // добавление текущого сокета с глобальной список сокетов
InputStream in = s.getInputStream(); ois = new ObjectInputStream(in);//получаем от игрока новое ребро
A = new int [5][5]; for (int i = 0; i<5; i++) for (int j = 0; j < 5; j++) { A[i][j] = 0; }
while (s.isConnected()) {
Para new_p = (Para)ois.readObject();
x = new_p.X(); y = new_p.Y(); p = new_p.P();
A[x][y] = p; System.out.println(x + ";" + y); ServerAnswer sa = new ServerAnswer(5, 5); sa.SetAns(new_p, A);
listSocket = ListSocket.getListSocket(); for (Socket socket : listSocket) { // отсылка сообщения всем сокетам
oos = new ObjectOutputStream(socket.getOutputStream()); oos.writeObject(sa); oos.flush(); }
} ListSocket.removeSocketWithList(s); // если поток завершается то сокет клиента удаляется из списка сокетов System.out.println("User disconnect..."); } catch (IOException ex) { try { s.close(); } catch (IOException ex1) { Logger.getLogger(SocketThread.class.getName()).log(Level.SEVERE, null, ex1); } Logger.getLogger(SocketThread.class.getName()).log(Level.SEVERE, null, ex); } catch (ClassNotFoundException ex) { Logger.getLogger(SocketThread.class.getName()).log(Level.SEVERE, null, ex); } } } Клиент: public class Client2 {
public static void main(String[] args) { try { System.out.println("Client starting..."); Socket s = new Socket("localhost",3129); System.out.println("Connect to server..."); Thread threadIn = new Thread(new SocketInputThread(s));// создание отдельного потока на считывание даных от сервера Thread threadOut = new Thread(new SocketOutputThread(s));// создание отдельного потока на ввод даных с клавиатуры threadIn.start(); threadOut.start(); } catch (UnknownHostException ex) { Logger.getLogger(Client2.class.getName()).log(Level.SEVERE, null, ex); } catch (IOException ex) { Logger.getLogger(Client2.class.getName()).log(Level.SEVERE, null, ex); } } } SocketInputThread: public class SocketInputThread implements Runnable {
private Socket s = null; private ObjectInputStream ois = null;
public SocketInputThread(Socket s) { this.s = s; }
@Override public void run() { try { InputStream in = s.getInputStream(); ois = new ObjectInputStream(in);
while(true){ ServerAnswer sa;
sa = (ServerAnswer)ois.readObject(); String str = "";
for (int i=0; iprivate Socket s = null; private ObjectOutputStream oos = null;
public SocketOutputThread(Socket s) { this.s = s; }
@Override public void run() { try { Scanner sc = new Scanner(System.in);
oos = new ObjectOutputStream(s.getOutputStream()); while (true) { int x = sc.nextInt(); int y = sc.nextInt(); Para p = new Para(); p.SetPara(x,y,1);
oos.writeObject(p); oos.flush(); } } catch (IOException ex) { Logger.getLogger(SocketOutputThread.class.getName()).log(Level.SEVERE, null, ex); } } }


Ответ

На один ObjectOutputStream должен приходиться ровно один ObjectInputStream Поясню: Корректный object stream выглядит так: [stream header], [object], [object], [object],... У вас object stream'ы выглядят так: [stream header], [object], [object], [stream header], [object], ... Stream header записывается каждый раз при создании нового ObjectOutputStream. Ваш invalid type code: AC означает, что ObjectInputStream ожидал прочитать объект, а прочитал хидер начала потока. То есть ошибка тут: for (Socket socket : listSocket) { // отсылка сообщения всем сокетам oos = new ObjectOutputStream(socket.getOutputStream()); Нужно хранить не список сокетов, чтобы каждый раз, проходясь по ним, создавать новые ObjectOutputStream'ы, а хранить список ObjectOutputStream'ов и переиспользовать их.

понедельник, 8 апреля 2019 г.

Что происходит когда сокет записывает данные в то время, когда с другой стороны читаются предыдущие данные?

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


Ответ

Если сокет блокирующий и на чтение запрошено больше данных, чем Вы отправляете - будет ждать "второй итерации". Если сокет неблокирующий - вычитает, сколько сможет, и затем Вам снова придется опрашивать сокет.
В связи с комментарием @VTT хочу дать немного более развернутый ответ.

Если вызывается одна из функций read/recv/recvfrom для блокируемого сокета и при этом в буфере нет никаких данных - сокет переходит в спящее состояние до тех пор, пока не придут какие-либо данные. Достаточно одного байта.
Тем не менее, мы можем задать флаг MSG_WAITALL - в этом случае мы будем ждать до тех пор, пока не будет доступно фиксированное (запрошенное нами) количество байт.
В случае неблокируемого сокета - если не удовлетворены условия ввода - то ф-ия вернет управление установив ошибку EWOULDBLOCK

Если вызывается одна из функций write/send/sendto - данные копируются из буфера приложения в буфер отправки сокета. Для блокируемого сокета - если в буфере отправки недостаточно места, процесс переходит в состоянии оидания до тех пор, пока это самое место не освободится.
Для неблокируемого сокета - если в буфере отправки недостаточно места - ф-ия вернет управление, установив ошибку EWOULDBLOCK

Комбинируйте и предполагайте наиболее вероятный результат того, чем же все-таки окончится посылка двух сообщений :)