【问题标题】:Netty client not receiving the complete data sent by the ServerNetty客户端没有收到Server发送的完整数据
【发布时间】:2018-11-28 19:20:04
【问题描述】:

我正在设计一个基于 Netty 的解决方案,通过 TCP 将文件从服务器传输到客户端。客户端指定文件的位置,然后服务器将文件发送给客户端。

目前,该解决方案适用于小文件(

如果要发送的文件大于 ~5MB,只发送部分数据,这会有所不同(每次发送的数据量都不相同)。此外,从日志中可以看出,服务器已经发送了完整的数据量(文件)。

问题是客户端没有收到服务器发送的完整数据。我下面的代码有什么问题?或有人能指出我正确的方向吗?

以下是我的客户端、服务器及其处理程序: (为了简洁我只列出了重要的方法)

客户:

 public class FileClient {

        private final static int PORT = 8992;
        private final static String HOST = "127.0.0.1";

        public class ClientChannelInitializer extends ChannelInitializer<SocketChannel> {

            private SslContext sslContext = null;
            private String srcFile = "";
            private String destFile = "";

            public ClientChannelInitializer(String srcFile, String destFile, SslContext sslCtx) {
                this.sslContext = sslCtx;
                this.srcFile = srcFile;
                this.destFile = destFile;
            }

            @Override
            protected void initChannel(SocketChannel socketChannel) throws Exception {
                ChannelPipeline pipeline = socketChannel.pipeline();
                pipeline.addLast(sslContext.newHandler(socketChannel.alloc(), HOST, PORT));
                pipeline.addLast("clientHandler", new FileClientHandler(srcFile, destFile));
            }

        }

        private void startUp(String srcFile, String destFile) throws Exception {
            SslContext sslCtx = SslContextBuilder.forClient().trustManager(InsecureTrustManagerFactory.INSTANCE).build();
            EventLoopGroup workerGroup = new NioEventLoopGroup();

                Bootstrap clientBootstrap = new Bootstrap();
                clientBootstrap.group(workerGroup);
                clientBootstrap.channel(NioSocketChannel.class);
                clientBootstrap.option(ChannelOption.TCP_NODELAY, true);
                clientBootstrap.handler(new LoggingHandler(LogLevel.INFO));
                clientBootstrap.handler(new ClientChannelInitializer(srcFile, destFile, sslCtx));

Channel channel = clientBootstrap.connect(new InetSocketAddress(HOST, PORT)).sync().channel();
                channel.closeFuture().sync();
            } 
        }

        public static void main(String[] args) throws Exception {

            String src = "/Users/home/src/test.mp4";
            String dest = "/Users/home/dest/test.mp4";
            new FileClient().startUp(src, dest);
        }

    }

客户端处理程序:

public class FileClientHandler extends SimpleChannelInboundHandler<ByteBuf> {


    private final String sourceFileName;
    private OutputStream outputStream;
    private Path destFilePath;
    private byte[] buffer = new byte[0];



    public FileClientHandler(String SrcFileName, String destFileName) {
        this.sourceFileName = SrcFileName;
        this.destFilePath = Paths.get(destFileName);
        System.out.println("DestFilePath-" + destFilePath);
    }

    @Override
    public void channelActive(ChannelHandlerContext ctx) throws Exception {
        ctx.writeAndFlush(ToByteBuff(this.sourceFileName));
    }

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, ByteBuf byteBuff) throws Exception {
        if (this.outputStream == null) {
            Files.createDirectories(this.destFilePath.getParent());
            if (Files.exists(this.destFilePath)) {
                Files.delete(this.destFilePath);
            }
            this.outputStream = Files.newOutputStream(this.destFilePath, StandardOpenOption.CREATE,
                    StandardOpenOption.APPEND);
        }

        int size = byteBuff.readableBytes();
        if (size > this.buffer.length) {
            this.buffer = new byte[size];
        }
        byteBuff.readBytes(this.buffer, 0, size);
        this.outputStream.write(this.buffer, 0, size);

    }   

文件服务器:

public class FileServer {
    private final int PORT = 8992;

    public void run() throws Exception {
        SelfSignedCertificate ssc = new SelfSignedCertificate();
        final SslContext sslCtx = SslContextBuilder.forServer(ssc.certificate(), ssc.privateKey()).build();
        EventLoopGroup bossGroup = new NioEventLoopGroup();
        EventLoopGroup workerGroup = new NioEventLoopGroup();
        try {
            ServerBootstrap b = new ServerBootstrap();
            b.group(bossGroup, workerGroup).channel(NioServerSocketChannel.class).option(ChannelOption.SO_BACKLOG, 100)
                    .handler(new LoggingHandler(LogLevel.INFO)).childHandler(new ChannelInitializer<SocketChannel>() {
                        @Override
                        public void initChannel(SocketChannel ch) throws Exception {
                            ChannelPipeline pipeline = ch.pipeline();
                            pipeline.addLast(sslCtx.newHandler(ch.alloc()));

                            pipeline.addLast(new ChunkedWriteHandler());
                            pipeline.addLast(new FilServerFileHandler());
                        }
                    });
            ChannelFuture f = b.bind(PORT).sync();

            f.channel().closeFuture().sync();
        } finally {
            bossGroup.shutdownGracefully();
            workerGroup.shutdownGracefully();
        }
    }

    public static void main(String[] args) throws Exception {
        new FileServer().run();
    }
}

文件服务器处理程序:

public class FilServerFileHandler extends SimpleChannelInboundHandler<ByteBuf> {

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, ByteBuf buff) throws Exception {
        String filePathStr = byteBuf.toString(CharsetUtil.UTF_8);

        File file = new File(filePathStr);
        RandomAccessFile raf = null;
        ChannelFuture sendFileFuture;
        try {
            raf = new RandomAccessFile(file, "r");

            sendFileFuture = ctx.writeAndFlush(new ChunkedNioFile(raf.getChannel()),
                    ctx.newProgressivePromise());

            sendFileFuture.addListener(new ChannelProgressiveFutureListener() {
                public void operationComplete(ChannelProgressiveFuture future) throws Exception {
                    System.err.println("Transfer complete.");
                }

                public void operationProgressed(ChannelProgressiveFuture future, long progress, long total)
                        throws Exception {
                    if (total < 0) { // total unknown
                        System.err.println("Transfer progress: " + progress);
                    } else {
                        System.err.println("Transfer progress: " + progress + " / " + total);
                    }
                }
            });

        } catch (FileNotFoundException fnfe) {
        } finally {
            if (raf != null)
                raf.close();
        }
    }

我检查了SO Q1 和SO Q2

【问题讨论】:

  • 为什么不使用这里的方法:stackoverflow.com/questions/25888260/… 给您的客户?
  • 谢谢。我也检查过这篇文章。即使这样也行不通。从服务器发送的任何数量 (> 5MB) 的数据中,最多只能接收/保存 5MB 的数据。
  • 尝试在您的客户端和服务器处理程序中添加一个 exceptionCaught 方法,并检查在传输过程中是否引发了异常。而且,您的服务器是否在任何类型的代理后面?像 nginx 或 haproxy
  • @kelgon 我确实有 exceptionCaught 方法,只是为了简洁起见,我在这里没有提到它。 1. 不,在传输过程中没有抛出异常。 2. 不,不涉及代理。谢谢。

标签: java tcp netty


【解决方案1】:

通过FilServerFileHandler 中的一些小调整解决了您的问题:

public class FileServerHandler extends SimpleChannelInboundHandler<ByteBuf> {
    @Override
    protected void channelRead0(ChannelHandlerContext ctx, ByteBuf buff) throws Exception {
        String filePathStr = buff.toString(CharsetUtil.UTF_8);

        File file = new File(filePathStr);
        RandomAccessFile raf = new RandomAccessFile(file, "r");
        ChannelFuture sendFileFuture;
        try {
            sendFileFuture = ctx.writeAndFlush(new ChunkedNioFile(raf.getChannel()), ctx.newProgressivePromise());
            sendFileFuture.addListener(new ChannelProgressiveFutureListener() {
                public void operationComplete(ChannelProgressiveFuture future) throws Exception {
                    System.err.println("Transfer complete.");
                    if (raf != null) {
                        raf.close();
                    }
                }
                public void operationProgressed(ChannelProgressiveFuture future, long progress, long total)
                        throws Exception {
                    if (total < 0) { // total unknown
                        System.err.println("Transfer progress: " + progress);
                    } else {
                        System.err.println("Transfer progress: " + progress + " / " + total);
                    }
                }
            });
        } catch (FileNotFoundException e) {
            e.printStackTrace();
        }
    }
}

我将raf.close() 移动到operationComplete 方法中。

部分传输是由于在写操作期间关闭raf引起的。注意ctx.writeAndFlush是一个异步调用,所以finally块中的raf.close()可能会在写操作完成之前被触发,尤其是当文件足够大的时候。

【讨论】:

  • 顺便说一句,如果您的操作系统支持零拷贝(即 Linux 中的 sendfile()),FileRegion 的性能比 ChunkedNioFile 更好,后者在文件系统和 nic 缓冲区之间直接传输字节,而 ChunkedNioFile 需要读取字节到用户空间(JVM 堆)
  • 好的。将纳入这一点。我不清楚一件事。使用TCP传输大文件时是否不需要使用某种LengthBasedDecoders?
  • 是的,它不是必需的。大数据将通过 TCP 以小数据包的形式发送,在 TCP 层进行分片和组装。客户端处理程序中的channelRead0 将被多次调用。
  • 另外我忘了说零拷贝传输文件数据,所以你不能使用ssl或任何类型的数据加密技术
  • kelgon 谢谢。好的,我知道了。只有在非 ssl 的情况下才会使用零拷贝。
猜你喜欢
  • 1970-01-01
  • 2013-03-14
  • 1970-01-01
  • 2015-09-28
  • 2015-02-10
  • 2018-01-30
  • 1970-01-01
  • 1970-01-01
  • 2014-12-27
相关资源
最近更新 更多