【问题标题】:ObjectInputStream returns ObjectStreamClassObjectInputStream 返回 ObjectStreamClass
【发布时间】:2015-01-13 00:09:56
【问题描述】:

我遇到了一些麻烦。以下行在第一次读取时在客户端上执行良好,但在第二次读取时失败。

WorkerNode.java:72 Message task = (Message) in.readObject();

in 是一个私有的 ObjectInputStream。收到的异常如下

Exception in thread "main" java.lang.ClassCastException: java.io.ObjectStreamClass cannot be cast to parallelprogramming.Message
    at parallelprogramming.WorkerNode.receiveTask(WorkerNode.java:72)
    at parallelprogramming.WorkerNode.computeTillEndOfWork(WorkerNode.java:139)
    at parallelprogramming.Worker.main(Worker.java:24)
Nov 15, 2014 11:07:15 PM parallelprogramming.WorkerNode receiveTask
SEVERE: null
java.io.StreamCorruptedException: invalid type code: 00
    at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1379)
    at java.io.ObjectInputStream.readObject(ObjectInputStream.java:371)
    at parallelprogramming.WorkerNode.receiveTask(WorkerNode.java:72)
    at parallelprogramming.WorkerNode.lambda$startListeningForWork$0(WorkerNode.java:59)
    at parallelprogramming.WorkerNode$$Lambda$1/798154996.run(Unknown Source)
    at java.lang.Thread.run(Thread.java:745)

Exception in thread "WorkListener" java.lang.ClassCastException: parallelprogramming.MatMulTask cannot be cast to parallelprogramming.Message
    at parallelprogramming.WorkerNode.receiveTask(WorkerNode.java:72)
    at parallelprogramming.WorkerNode.lambda$startListeningForWork$0(WorkerNode.java:59)
    at parallelprogramming.WorkerNode$$Lambda$1/798154996.run(Unknown Source)
    at java.lang.Thread.run(Thread.java:745)
Java Result: 1

消息具有以下结构。

public class Message implements IMessage{
    private final MessageType type;
    private final Object payload;

IMessage 扩展了 Serializable。我在客户端和服务器端都使用相同的 ObjectInputStream 和 ObjectOutputStream。我试着四处寻找,但没有运气。其他人发现了类似的东西吗?

编辑2: 向 Worker 发送消息的代码:

private final Map<String, WorkerConn> nodes;
//nodes initialized in constructor
private void sendTaskToNode(ITask task, String node) {
        if(task == null){
            return;
        }        
        try{   
            Message msg = new Message(MessageType.task, task);
            nodes.get(node).sendObject(msg);

            nodes.get(node).incrementworkCount();

            System.out.println("Sent work to "+node);
        } catch (IOException ex) {
            Logger.getLogger(MasterNode.class.getName()).log(Level.SEVERE, null, ex);
        }
    }

在 WorkerConn 中发送对象:

public WorkerConn(Socket socket, String name) throws IOException {
        this.socket = socket;
        this.out = new ObjectOutputStream(socket.getOutputStream());
        this.in = new ObjectInputStream(socket.getInputStream());
        this.name = name;
        workCount = 0;
    }

void sendObject(Message msg) throws IOException {
        out.writeObject(msg);
    }

在Worker上接收Message的部分:

public void receiveTask() throws NotConnectedToMasterException {
        try {
            Message task = (Message) in.readObject();
            if(task.getMessageType() == MessageType.task){
                tasks.add((ITask) task.getPayload());
                System.out.println("Received Task");
            }else if(task.getMessageType() == MessageType.endOfWork){
                ITask t = new AbstractTask() {

                    @Override
                    public Object call() throws Exception {
                        throw new UnsupportedOperationException("Not supported yet."); //To change body of generated methods, choose Tools | Templates.
                    }
                };
                t.setDeathPill();
                tasks.add(t);
                System.out.println("added deathpill");
            }else{
                System.out.println("Received "+task.getMessageType());
            }
        } catch(EOFException ex) {
            return;
        } catch (IOException ex) {
            Logger.getLogger(WorkerNode.class.getName()).log(Level.SEVERE, null, ex);
            return;
        } catch (ClassNotFoundException ex) {
            Logger.getLogger(WorkerNode.class.getName()).log(Level.SEVERE, null, ex);
        }
    }

编辑: 发现了我的问题。我正在创建一个线程来侦听传入任务,读取输入流,并且在任务队列为空时还从主线程读取相同的输入流。在 Worker 节点中:

private void startListeningForWork(){
        workListener = new Thread(() -> {
            while(!master.isClosed()){
                try {
                    receiveTask();
                } catch (NotConnectedToMasterException ex) {
                    break;
                }
            }
        });
        workListener.setName("WorkListener");
        workListener.start();
    }

while(!task.isDeathPill()){
            try {                
                results.addToResult(task.call());
                sendACKtoMaster();                
                task = tasks.remove();
            } catch (Exception ex) {
                try {
                    receiveTask();
                } catch (NotConnectedToMasterException ex1) {
                    Logger.getLogger(WorkerNode.class.getName()).log(Level.SEVERE, null, ex1);
                    break;
                }
            }   
        }

这导致其中一个线程先于另一个线程读取对象并搞砸了。

【问题讨论】:

  • 同时显示您发送数据的位置
  • 已将整个代码上传到 github。文章末尾的链接
  • 很抱歉,但我不会在所有这些文件中查看与特定 readObject 声明相关的单个 writeObject 声明(尤其是因为我正在打电话)。至少告诉我哪个类在哪个包中
  • 对此感到抱歉。消息在这里发送github.com/kalgecin/parallelProgrammingTest/blob/master/src/…
  • 没有详细介绍,但我很确定您有一些同步问题。 WorkerNode 以线程不安全的方式使用节点映射。我没有完全分析工作人员,但我的猜测是您没有正确同步对流的访问。

标签: java serialization stream network-programming


【解决方案1】:

发现我的问题。我正在创建一个线程来侦听传入任务,读取输入流,并且在任务队列为空时还从主线程读取相同的输入流。在 Worker 节点中:

private void startListeningForWork(){
        workListener = new Thread(() -> {
            while(!master.isClosed()){
                try {
                    receiveTask();
                } catch (NotConnectedToMasterException ex) {
                    break;
                }
            }
        });
        workListener.setName("WorkListener");
        workListener.start();
    }

while(!task.isDeathPill()){
            try {                
                results.addToResult(task.call());
                sendACKtoMaster();                
                task = tasks.remove();
            } catch (Exception ex) {
                try {
                    receiveTask();
                } catch (NotConnectedToMasterException ex1) {
                    Logger.getLogger(WorkerNode.class.getName()).log(Level.SEVERE, null, ex1);
                    break;
                }
            }   
        }

这导致其中一个线程先于另一个线程读取对象并搞砸了。

【讨论】:

  • 正如@jtahlborn 预测的那样。
  • 是的 :) 在我发现它之前调试了每一行:( java 需要更好地解释这种东西的异常:p 或更好地处理比赛
猜你喜欢
  • 1970-01-01
  • 2017-07-17
  • 2018-08-25
  • 2011-09-02
  • 1970-01-01
  • 2020-05-13
  • 1970-01-01
  • 2012-07-29
  • 1970-01-01
相关资源
最近更新 更多