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

Поиск:

Ответ в темуСоздание новой темы Создание опроса
> Java Server c BlockingQueue, Java Server c BlockingQueue 
:(
    Опции темы
tengorn
Дата 24.9.2011, 15:13 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Новичок



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

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



доброе время суток!

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

вот весь пакет для тестов

Сервер
Код

import java.net.DatagramPacket;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;

public class Server {

    private static int port = 12345;
    
    private static BlockingQueue<DatagramPacket> queue;

    public static void main(String[] args) {

        queue = new LinkedBlockingQueue<DatagramPacket>();

            ServerProducer producer = new ServerProducer(queue,port);
            ServerConsumer consumer = new ServerConsumer(queue);

            new Thread(producer).start();
            new Thread(consumer).start();
    }

}



Consumer
Код

import java.io.ByteArrayInputStream;
import java.io.ObjectInputStream;
import java.net.DatagramPacket;
import java.util.concurrent.BlockingQueue;

public class ServerConsumer implements Runnable {

    protected BlockingQueue<DatagramPacket> queue = null;

    private ObjectInputStream ois;

    private Package receivedPackage;
    
    private DatagramPacket received;

    public ServerConsumer(BlockingQueue<DatagramPacket> queue) {
        this.queue = queue;
    }

    public void run() {

        try {
            while (true) {

                //System.out.println(queue.take());

                received = queue.take();
                ByteArrayInputStream bais = new ByteArrayInputStream(received.getData());
                ois = new ObjectInputStream(bais);
                receivedPackage = (Package) ois.readObject();
                System.out.println(receivedPackage.getName());
                System.out.println(receivedPackage.getNr());

            }
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}


Producer 
Код

import java.net.DatagramPacket;
import java.net.DatagramSocket;
import java.util.concurrent.BlockingQueue;


public class ServerProducer implements Runnable {

    private DatagramSocket serverSocket;

    private final byte[] receiveData = new byte[256];

    private DatagramPacket receivePackage;

    private int port;

    private BlockingQueue<DatagramPacket> queue;

    public ServerProducer(BlockingQueue<DatagramPacket> queue, int port) {
        this.queue = queue;
        this.port = port;
        try {
            serverSocket = new DatagramSocket(this.port);
            receivePackage = new DatagramPacket(receiveData, receiveData.length);
        } catch (Exception e) {
        }

        System.out.println("Server starts...");
    }

    public void run() {

        int i =1;
        try {
            while (true) {

                serverSocket.receive(receivePackage);
                queue.put(receivePackage);
                System.out.println("getted "+i);
                i++;

            }
        } catch (Exception e) {
        }

    }
}


Посылаемый обьект
Код

import java.io.Serializable;

public class Package implements Serializable {

    private static final long serialVersionUID = 1L;
    private String name;
    private int nr;


    public Package(String name, int nr) {
        this.name = name;
        this.nr = nr;
    }

    public String getName() {
        return name;
    }

    public int getNr() {
        return nr;
    }
}



Клиент
Код

import java.io.ByteArrayOutputStream;
import java.io.ObjectOutputStream;
import java.net.DatagramPacket;
import java.net.DatagramSocket;
import java.net.InetAddress;
import java.util.Vector;

public class Client {

    private ObjectOutputStream out = null;
    private InetAddress serverIP = null;
    private int port = 12345;

    private final byte[] ip = new byte[] { (byte) 127, (byte) 0, (byte) 0,
            (byte) 1 };

    public Client() throws Exception {

        final DatagramSocket clientSocket = new DatagramSocket();
        serverIP = InetAddress.getByAddress(ip);

        Vector<Package> vt = new Vector<Package>();
        vt.add(new Package("test1", 1));
        vt.add(new Package("test2", 2));
        vt.add(new Package("test3", 3));
        vt.add(new Package("test4", 4));
        vt.add(new Package("test5", 5));

        for (int i = 0; i < vt.size(); i++) {

            ByteArrayOutputStream baos = new ByteArrayOutputStream();
            out = new ObjectOutputStream(baos);
            out.writeObject(vt.get(i));
            out.flush();

            byte[] sendData = baos.toByteArray();

            DatagramPacket sendPacket = new DatagramPacket(sendData,
                    sendData.length, serverIP, port);
            clientSocket.send(sendPacket);
        }
    }

    public static void main(String[] args) {
        try {
            new Client();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }

}



проблема в том что если просто System.out.println(queue.take()); в ServerConsumer выдавать без извлечения обьекта, то выдаются все полученные пакеты
как только я хочу извлечь обьект и выдать его данные на консоль, то выдаётся только последний обьект и больше ничего :(

помогите устранить ошибку smile

PM MAIL   Вверх
jpr111a
Дата 10.10.2011, 20:14 (ссылка) | (нет голосов) Загрузка ... Загрузка ... Быстрая цитата Цитата


Новичок



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

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



Если ещё интересно, то проблема решается переносом
receivePackage = new DatagramPacket(receiveData, receiveData.length);
внутрь цикла в ServerProducer. При теперешней имплементции в очереди оказываются одинаковые объекты. 
В результате ServerProducer ждёт следующую датаграмму serverSocket.receive(receivePackage) захватив монитор (см. исходники метода). А ServerConsumer ждёт освобождения монитора received.getData()
PS. Комментарии о том, что в этом коде является плохим для проакшена, опускаю.

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

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

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


 




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


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

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