Версия для печати темы
Нажмите сюда для просмотра этой темы в оригинальном формате
Форум программистов > Java: Работа с сетью > Java Server c BlockingQueue


Автор: tengorn 24.9.2011, 15:13
доброе время суток!

хочеться сделать сервак через 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

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

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