【问题标题】:Java: Asynchronous I/O channel for reading and writing linesJava:用于读写行的异步 I/O 通道
【发布时间】:2016-09-29 10:31:51
【问题描述】:

我有一个应用程序,它使用 BufferedReaderPrintStream 包装 InputStreamOutputStream 对象的 OutputStream 同步读取和写入文本行。所以,我可以只使用BufferedReader.readLine()PrintStream.println() 方法,让Java 库将输入分成几行并为我格式化输出。

现在我想用异步 IO 替换这个同步 IO。所以我一直在研究AsynchronousSocketChannel,它允许异步读写字节。现在,我想要包装类,以便我可以使用字符串异步读取/写入行。

我在 Java 库中找不到这样的包装类。在我编写自己的实现之前,我想问一下是否有任何其他库允许包装 AsynchronousSocketChannel 并提供异步文本 IO。

【问题讨论】:

  • 为什么?你想解决什么问题?使用BufferedReader,您可以每秒读取数百万行。这还不够吗?
  • @EJP:我想异步阅读。我不想阻止等待通过套接字接收一行文本。我希望在收到完整的文本行后调用我的代码。
  • @giorgio-b 如果你不是从套接字读取,那会收到完整的行吗?
  • @EJP:我的问题很笼统。如果我有同步字节 IO 和同步文本 IO 的各种包装器,我希望异步 IO 有相同的。
  • 我就是这么说的。它不在那里。场外资源问题不在此处讨论。

标签: java asynchronous text io


【解决方案1】:

你可以这样做

public void nioAsyncParse(AsynchronousSocketChannel channel, final int bufferSize) throws IOException, ParseException, InterruptedException {
    ByteBuffer byteBuffer = ByteBuffer.allocate(bufferSize);
    BufferConsumer consumer = new BufferConsumer(byteBuffer, bufferSize);
    channel.read(consumer.buffer(), 0l, channel, consumer);
}


class BufferConsumer implements CompletionHandler<Integer, AsynchronousSocketChannel> {

        private ByteBuffer bytes;
        private StringBuffer chars;
        private int limit;
        private long position;
        private DateFormat frmt = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");

        public BufferConsumer(ByteBuffer byteBuffer, int bufferSize) {
            bytes = byteBuffer;
            chars = new StringBuffer(bufferSize);
            frmt = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
            limit = bufferSize;
            position = 0l;
        }

        public ByteBuffer buffer() {
            return bytes;
        }

        @Override
        public synchronized void completed(Integer result, AsynchronousSocketChannel channel) {

            if (result!=-1) {
                bytes.flip();
                final int len = bytes.limit();
                int i = 0;
                try {
                    for (i = 0; i < len; i++) {
                        byte by = bytes.get();
                        if (by=='\n') {
                            // ***
                            // The code used to process the line goes here
                            // ***
                            chars.setLength(0);
                        }
                        else {
                            chars.append((char) by);
                        }
                    }
                }
                catch (Exception x) {
                    System.out.println("Caught exception " + x.getClass().getName() + " " + x.getMessage() + " i=" + String.valueOf(i) + ", limit=" + String.valueOf(len) + ", position="+String.valueOf(position));
                }

                if (len==limit) {
                    bytes.clear();
                    position += len;
                    channel.read(bytes, position, channel, this);
                }
                else {
                    try {
                        channel.close();
                    }
                    catch (IOException e) { }
                    bytes.clear();
                    buffers.add(bytes);
                }
            }
            else {
                try {
                    channel.close();
                }
                catch (IOException e) { }
                bytes.clear();
                buffers.add(bytes);
            }
        }

        @Override
        public void failed(Throwable e, AsynchronousSocketChannel channel) {
        }
};

【讨论】:

    猜你喜欢
    • 2013-06-19
    • 2012-09-14
    • 1970-01-01
    • 1970-01-01
    • 2015-07-24
    • 1970-01-01
    • 2011-03-14
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多