【问题标题】:Java Socket InputStream read missing bytesJava Socket InputStream 读取丢失的字节
【发布时间】:2015-07-11 05:25:45
【问题描述】:

通过从套接字的输入流中读取字节,我得到了一个非常奇怪的行为。

在我的项目中,客户向服务发出请求。对于每个请求,都会建立一个新连接。

首先发送的字节告诉服务将遵循什么样的请求。

然后发送请求本身。

服务接收字节并继续请求。这确实适用于所有请求的至少 95%。剩下的 5% 有一个奇怪的行为,我想不通。

字节并不是发送的所有字节。但是关于这个话题最奇怪的事情是丢失的字节不在流的开头或结尾。它们遍布整个溪流。

很遗憾,我无法在此处提供完整代码,因为它与工作相关。但我可以提供显示问题本身的测试代码。

为了弄清楚发生了什么,我写了 2 个类。一个来自java.net.Socket,另一个来自java.net.ServerSocket。

代码如下:

import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.net.Socket;
import java.util.ArrayList;
import java.util.List;



public class DebugSocket extends Socket
{
    private class InputStreamWrapper extends InputStream
    {
        private int
            availables,
            closes,
            marksupporteds,
            resets;

        private List<Integer>
            marks   = new ArrayList<Integer>(),
            reads   = new ArrayList<Integer>();

        private List<Long>
            skips   = new ArrayList<Long>();


        @Override
        public int available() throws IOException
        {
            availables++;
            return DebugSocket.this.origininput.available();
        }

        @Override
        public void close() throws IOException
        {
            closes++;
            DebugSocket.this.origininput.close();
        }

        @Override
        public synchronized void mark(int readlimit)
        {
            marks.add(readlimit);
            DebugSocket.this.origininput.mark(readlimit);
        }

        @Override
        public boolean markSupported()
        {
            marksupporteds++;
            return DebugSocket.this.origininput.markSupported();
        }

        @Override
        public synchronized void reset() throws IOException
        {
            resets++;
            DebugSocket.this.origininput.reset();
        }

        @Override
        public int read() throws IOException
        {
            int read = DebugSocket.this.origininput.read();

            reads.add(read);

            if ( read != -1 )
            {
                DebugSocket.this.inputdebugbuffer.write(read);
            }

            return read;
        }

        @Override
        public int read(byte[] b) throws IOException
        {
            int read = DebugSocket.this.origininput.read(b);

            DebugSocket.this.inputdebugbuffer.write(b, 0, read);

            return read;
        }

        @Override
        public int read(byte[] b, int off, int len) throws IOException
        {
            int read = DebugSocket.this.origininput.read(b, off, len);

            DebugSocket.this.inputdebugbuffer.write(b, off, read);

            return read;
        }

        @Override
        public long skip(long n) throws IOException
        {
            long skipped = DebugSocket.this.origininput.skip(n);

            skips.add(skipped);

            return skipped;
        }
    }

    private class OutputStreamWrapper extends OutputStream
    {
        private int
            flushes,
            closes;


        @Override
        public void close() throws IOException
        {
            closes++;
            DebugSocket.this.originoutput.close();
        }

        @Override
        public void flush() throws IOException
        {
            flushes++;
            DebugSocket.this.originoutput.flush();
        }

        @Override
        public void write(int b) throws IOException
        {
            DebugSocket.this.outputdebugbuffer.write(b);
            DebugSocket.this.originoutput.write(b);
            DebugSocket.this.originoutput.flush();
        }

        @Override
        public void write(byte[] b) throws IOException
        {
            DebugSocket.this.outputdebugbuffer.write(b);
            DebugSocket.this.originoutput.write(b);
            DebugSocket.this.originoutput.flush();
        }

        @Override
        public void write(byte[] b, int off, int len) throws IOException
        {
            DebugSocket.this.outputdebugbuffer.write(b, off, len);
            DebugSocket.this.originoutput.write(b, off, len);
            DebugSocket.this.originoutput.flush();
        }
    }


    private static final Object
        staticsynch = new Object();

    private static long
        idcounter   = 0;


    private final long
        id;

    private final ByteArrayOutputStream
        inputdebugbuffer,
        outputdebugbuffer;

    private final InputStream
        inputwrapper;

    private final OutputStream
        outputwrapper;

    private InputStream
        origininput;

    private OutputStream
        originoutput;


    public InputStream getInputStream() throws IOException
    {
        if ( origininput == null )
        {
            synchronized ( inputdebugbuffer )
            {
                if ( origininput == null )
                {
                    origininput = super.getInputStream();
                }
            }
        }

        return inputwrapper;
    }

    public OutputStream getOutputStream() throws IOException
    {
        if ( originoutput == null )
        {
            synchronized ( outputdebugbuffer )
            {
                if ( originoutput == null )
                {
                    originoutput    = super.getOutputStream();
                }
            }
        }

        return outputwrapper;
    }


    public DebugSocket()
    {
        id                  = getNextId();
        inputwrapper        = new InputStreamWrapper();
        outputwrapper       = new OutputStreamWrapper();
        inputdebugbuffer    = new ByteArrayOutputStream();
        outputdebugbuffer   = new ByteArrayOutputStream();
    }


    private static long getNextId()
    {
        synchronized ( staticsynch )
        {
            return ++idcounter;
        }
    }
}
import java.io.IOException;
import java.net.ServerSocket;


public class DebugServerSocket extends ServerSocket
{
    public DebugServerSocket() throws IOException
    {
        super();
    }

    public DebugSocket accept() throws IOException
    {
        DebugSocket s = new DebugSocket();

        implAccept(s);

        return s;
    }
}

DebugSocket 类接收与InputStream 和OutputStream 的每次交互的通知

现在当问题发生时,我总是可以看到字节丢失。

这里是一个例子:

客户端发送 1758 个字节。我从DebugSocket 中的成员outputdebugbuffer 获得了前23 个字节。

Bytes: 0,0,0,0,0,0,0,2,0,0,6,-46,31,-117,8,0,0,0,0,0,0,0,-83

服务器收到 227 个字节。对于调试问题,我总是会读取输入流,直到我得到 -1,以便所有字节都得到处理。现在,我从 DebugSocket 中的成员 inputdebugbuffer 获得的服务器端的 16 个前导字节。

Bytes: 0,0,0,6,-46,31,-117,8,0,0,0,0,0,0,0,-83

如图所示,缺少 7 个字节。前 8 个字节是一个长值,我将其更改为一个字节值以进行调试。所以我认为第一个字节总是正确的。

如果是代码失败,则不会进行任何请求,但正如我之前所说,这种情况最多只发生在所有连接的 5% 上。

有人知道这里发生了什么吗?

我还使用DataInputStream 和DataOutputStream 发送数据。正如您在DebugSocket 的OutputStreamWrapper 中看到的那样,我总是在每次写入操作后刷新。

我错过了什么吗?

如果需要其他代码,我会尝试发布它。

附:该服务是多线程的,可以并行处理 100 个请求。此外,客户端是多线程的,并行执行 20 个请求。如前所述,每个请求都使用其一个连接,并在请求进行后立即关闭该连接。

我希望有人对这个问题有所了解。

编辑:

没有主要的方法可以显示它在 cmets 中所做的任何事情,但这里是使用的客户端和服务器的代码块。

客户端:(以 20 个线程并行运行)

    public void sendRequest(long _requesttype, byte[] _bytes)
    {
        Socket              socket  = null;
        DataInputStream     input   = null;
        DataOutputStream    output  = null;
        InputStream         sinput  = null;
        OutputStream        soutput = null;

        try
        {
            socket  = new DebugSocket();

            socket.connect(serveraddress);

            sinput  = socket.getInputStream();
            soutput = socket.getOutputStream();

            input   = new DataInputStream(sinput);
            output  = new DataOutputStream(soutput);

            output.writeLong(_requesttype);

            output.flush();
            soutput.flush();

            output.write(_bytes);

            output.flush();
            soutput.flush();

            // wait for notification byte that service had received all data.
            input.readByte();
        }
        catch (IOException ex)
        {
            LogHelper.log(ex);
        }
        catch (Error err)
        {
            throw err;
        }
        finally
        {
            output.flush();
            soutput.flush();

            input.close();
            output.close();

            finishSocket(socket);
        }
    }

服务器:(每个请求在一个线程中运行。最多 100 个线程)

    public void proceedRequest(DebugSocket _socket)
    {
        DataInputStream     input   = null;
        DataOutputStream    output  = null;
        InputStream         sinput  = null;
        OutputStream        soutput = null;

        try
        {
            sinput  = _socket.getInputStream();
            soutput = _socket.getOutputStream();

            input   = new DataInputStream(sinput);
            output  = new DataOutputStream(soutput);

            RequestHelper.proceed(input.readLong(), input, output);

            // send notification byte to the client.
            output.writeByte(1);

            output.flush();
            soutput.flush();
        }
        catch (IOException ex)
        {
            LogHelper.log(ex);
        }
        catch (Error err)
        {
            throw err;
        }
        finally
        {
            output.flush();
            soutput.flush();

            input.close();
            output.close();
        }
    }

在服务器代码中,readLong() 已经失败,原因是缺少字节。

【问题讨论】:

  • 您发布了一些代码,但我不确定您发布的代码是否与您的问题相关 - 也许它与您的调试相关,但在这种情况下,您可以说“我收到字节 [...]" 并没有准确解释你用来找出它的代码。
  • 如果您可以将问题缩小到minimal, complete and verifiable example,那会非常有用。
  • 你就在那里,但我试图避免在与工作相关的代码中出现关于我的读写行为的问题。使用这 2 个类并在最低级别显示问题。用 DebugSocket 完成什么操作并不重要,它都会被记录下来。
  • 不过,您的代码应该是完整的。应该有一个main,如果这是客户端和服务器,可能有两个。应该有显示问题的示例输入。应该有一些显示问题的东西,这样我们就可以在运行它时检查我们是否遇到了同样的问题。关键是我们能够重现您的问题,然后我们可以提供帮助。
  • 根据您的字节“值”,我猜您没有正确读取流。可能使用了错误的编码。所有对字节的读取都应返回 0 到 255 之间的值,包括 0 到 255。你不会得到负值,除了 -1 流关闭时。

标签: java sockets


【解决方案1】:

我的猜测是代码中的其他地方存在错误。我从问题中复制了DebugSocket 类并创建了一个MCVE(见下文)。它工作正常,我无法重现“服务器无法读取长值”问题。尝试修改下面的代码以包含更多您自己的代码,直到您可以重现问题为止,这应该让您知道在哪里寻找根本原因。

import java.io.*;
import java.net.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.*;

public class TestDebugSocket implements Runnable, Closeable {

    public static void main(String[] args) {

        TestDebugSocket m = new TestDebugSocket();
        try {
            m.run();
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            m.close();
        }
    }

    final int clients = 20;
    final boolean useDebugSocket = true;
    final byte[] someBytes = new byte[1758];
    final ThreadPoolExecutor tp = (ThreadPoolExecutor) Executors.newCachedThreadPool();
    final AtomicLong clientId = new AtomicLong();
    final ConcurrentLinkedQueue<Closeable> closeables = new ConcurrentLinkedQueue<Closeable>();
    final long maxWait = 5_000L;
    final CountDownLatch serversReady = new CountDownLatch(clients);
    final CountDownLatch clientsDone = new CountDownLatch(clients);

    ServerSocket ss;
    int port;

    @Override public void run()  {

        try {
            ss = useDebugSocket ? new DebugServerSocket() : new ServerSocket();
            ss.bind(null);
            port = ss.getLocalPort();
            tp.execute(new SocketAccept());
            for (int i = 0; i < clients; i++) {
                ClientSideSocket css = new ClientSideSocket();
                closeables.add(css);
                tp.execute(css);
            }
            if (!clientsDone.await(maxWait, TimeUnit.MILLISECONDS)) {
                System.out.println("CLIENTS DID NOT FINISH");
            } else {
                System.out.println("Finished");
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            close();
        }
    } 

    @Override public void close() {

        try { if (ss != null) ss.close(); } catch (Exception ignored) {}
        Closeable c = null;
        while ((c = closeables.poll()) != null) {
            try { c.close(); } catch (Exception ignored) {}
        }
        tp.shutdownNow();
    }

    class DebugServerSocket extends ServerSocket {

        public DebugServerSocket() throws IOException {
            super();
        }

        @Override public DebugSocket accept() throws IOException {

            DebugSocket s = new DebugSocket();
            implAccept(s);
            return s;
        }
    }

    class SocketAccept implements Runnable {

        @Override public void run() {
            try {
                for (int i = 0; i < clients; i++) {
                    SeverSideSocket sss = new SeverSideSocket(ss.accept());
                    closeables.add(sss);
                    tp.execute(sss);
                }
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    }

    class SeverSideSocket implements Runnable, Closeable {

        Socket s;
        public SeverSideSocket(Socket s) {
            this.s = s;
        }
        @Override public void run() {

            Long l = -1L;
            byte[] received = new byte[someBytes.length];
            try {
                DataInputStream in = new DataInputStream(s.getInputStream());
                DataOutputStream out = new DataOutputStream(s.getOutputStream());
                serversReady.countDown();
                if (!serversReady.await(maxWait, TimeUnit.MILLISECONDS)) {
                    System.out.println("CLIENTS DID NOT CONNECT ON TIME TO SERVER");
                }
                l = in.readLong();
                in.readFully(received);
                out.writeByte(1);
                out.flush();
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                // write to console at end to prevent synchronized socket I/O
                System.out.println("received long: " + l);
                close();
            }
        }

        @Override public void close() {
            TestDebugSocket.close(s);
            s = null;
        }
    }

    class ClientSideSocket implements Runnable, Closeable {

        Socket s;

        @SuppressWarnings("resource")
        @Override public void run() {

            Long l = -1L;
            Byte b = -1;
            try {
                s = useDebugSocket ? new DebugSocket() : new Socket();
                s.connect(new InetSocketAddress(port));
                DataInputStream in = new DataInputStream(s.getInputStream());
                DataOutputStream out = new DataOutputStream(s.getOutputStream());
                l = clientId.incrementAndGet();
                out.writeLong(l);
                out.write(someBytes);
                out.flush();
                b = in.readByte();
            } catch (Exception e) {
                e.printStackTrace();
            } finally {
                System.out.println("long send: " + l + ", result: " + b);
                close();
                clientsDone.countDown();
            }
        }

        @Override public void close() {
            TestDebugSocket.close(s);
            s = null;
        }
    }

    static void close(Socket s) {

        try { if (s != null) s.close(); } catch (Exception ignored) {}
    }
}

【讨论】:

  • 嗯,是的,我也是这样做的,如果您的系统输出匹配,那么我会遇到一个有趣的错误,因为我无法弄清楚可能会弄乱字节顺序的原因。我已经在 DebugSocket 中添加了一个记录,它使用 CallStack 以及线程。这里也没有失败。只有一个线程与套接字交互,并且调用仅来自众所周知和预期的位置。我会进一步研究它,如果我找到原因我会发布它。
【解决方案2】:

好的,我已经用所有可能的方法来定位原因。根据我在套接字编程和并行处理方面的经验,我可以说代码本身没有错误。嗅探器也告诉我。我的机器上有东西干扰了传输。

我停用了所有我能想到的(防火墙/防病毒/恶意软件扫描程序),但没有任何效果。

有人知道 tcp 包还有什么问题?

编辑:

好的,我明白了。 AVG 2014 搞砸了。 Jetzt 停用组件不起作用。在 Options->Settings 中有一个菜单点,您可以停用 AVG-Protection。

有人知道这方面的知识吗?

【讨论】:

  • 报告了同样的问题here。使用MSE 而不是 AVG(我使用 MSE 但没有问题)?
  • 哇,谢谢。这有助于解决问题。问题是 5 天前我有一个测试运行了将近 4 小时,没有任何问题,然后它就在那里。现在我只需要几秒钟就可以弹出问题^^
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2011-08-07
  • 2015-03-03
  • 2017-07-15
  • 1970-01-01
  • 2013-08-14
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多