【问题标题】:Read the newly appended file content to an InputStream in Java将新附加的文件内容读取到 Java 中的 InputStream
【发布时间】:2016-09-28 07:07:35
【问题描述】:

我有一个编写程序,它以特定速度将一个巨大的序列化 java 对象(以 1GB 为规模)写入本地磁盘上的二进制文件。实际上,writer 程序(用 C 语言实现)是一个网络接收器,它从远程服务器接收序列化对象的字节。 writer 的实现是固定的。

现在,我想实现一个 Java 阅读器程序,它可以读取文件并将其反序列化为 Java 对象。由于文件可能非常大,因此减少反序列化对象的延迟是有益的。特别是,我希望 Java 阅读器在对象的第一个字节被写入磁盘文件后开始读取/反序列化对象,这样即使在整个序列化对象被写入文件之前,阅读器也可以开始反序列化对象.阅读器提前知道文件的大小(在第一个字节写入文件之前)。

我认为我需要的是类似于阻塞文件 InputStream 的东西,当它到达 EndOfFile 时会被阻塞,但它没有读取预期的字节数(文件的大小将是)。因此,每当新字节写入文件时,读取器的 InputStream 就可以继续读取新内容。但是,Java 中的 FileInputStream 不支持此功能。

可能,我还需要一个文件侦听器来监视对文件所做的更改以实现此功能。

我想知道是否有任何现有的解决方案/库/包可以实现此功能。可能这个问题可能与监控日志文件中的一些问题类似。

字节流是这样的: FileInputStream -> SequenceInputStream -> BufferedInputStream -> JavaSerializer

【问题讨论】:

  • 这是一个 java nio Pipe 将输出连接到输入的问题。你当然需要线程,因此管道从未流行过。搜索示例。
  • 其他问题。您可以使用文件通道和内存映射字节缓冲区进行传输。如果你有一个巨大的序列化,检查所有的内部类可能是静态的,所以外部的 this 没有序列化。

标签: java logging serialization deserialization inputstream


【解决方案1】:

您需要两个线程:线程 1 从服务器下载并写入文件,线程 2 在文件可用时读取文件。

两个线程应该共享一个 RandomAccessFile,因此可以正确同步对 OS 文件的访问。你可以使用这样的包装类:

public class ReadWriteFile {
    ReadWriteFile(File f, long size) throws IOException {
        _raf = new RandomAccessFile(f, "rw");
        _size = size;

        _writer = new OutputStream() {

            @Override
            public void write(int b) throws IOException {
                write(new byte[] {
                        (byte)b
                });
            }

            @Override
            public void write(byte[] b, int off, int len) throws IOException {
                if (len < 0)
                    throw new IllegalArgumentException();
                synchronized (_raf) {
                    _raf.seek(_nw);
                    _raf.write(b, off, len);
                    _nw += len;
                    _raf.notify();
                }
            }
        };
    }

    void close() throws IOException {
        _raf.close();
    }

    InputStream reader() {
        return new InputStream() {
            @Override
            public int read() throws IOException {
                if (_pos >= _size)
                    return -1;
                byte[] b = new byte[1];
                if (read(b, 0, 1) != 1)
                    throw new IOException();
                return b[0] & 255;
            }

            @Override
            public int read(byte[] buff, int off, int len) throws IOException {
                synchronized (_raf) {
                    while (true) {
                        if (_pos >= _size)
                            return -1;
                        if (_pos >= _nw) {
                            try {
                                _raf.wait();
                                continue;
                            } catch (InterruptedException ex) {
                                throw new IOException(ex);
                            }
                        }
                        _raf.seek(_pos);
                        len = (int)Math.min(len, _nw - _pos);
                        int nr = _raf.read(buff, off, len);
                        _pos += Math.max(0, nr);
                        return nr;
                    }
                }
            }

            private long _pos;
        };
    }

    OutputStream writer() {
        return _writer;
    }

    private final RandomAccessFile _raf;
    private final long _size;
    private final OutputStream _writer;
    private long _nw;
}

以下代码展示了如何从两个线程中使用 ReadWriteFile:

public static void main(String[] args) throws Exception {
    File f = new File("test.bin");
    final long size = 1024;
    final ReadWriteFile rwf = new ReadWriteFile(f, size);

    Thread t1 = new Thread("Writer") {
        public void run() {
            try {
                OutputStream w = new BufferedOutputStream(rwf.writer(), 16);
                for (int i = 0; i < size; i++) {
                    w.write(i);
                    sleep(1);
                }
                System.out.println("Write done");
                w.close();
            } catch (Exception ex) {
                ex.printStackTrace();
            }
        }
    };

    Thread t2 = new Thread("Reader") {
        public void run() {
            try {
                InputStream r = new BufferedInputStream(rwf.reader(), 13);
                for (int i = 0; i < size; i++) {
                    int b = r.read();
                    assert (b == (i & 255));
                }
                int eof = r.read();
                assert (eof == -1);
                r.close();
                System.out.println("Read done");
            } catch (IOException ex) {
                ex.printStackTrace();
            }
        }
    };

    t1.start();
    t2.start();
    t1.join();
    t2.join();
    rwf.close();
}

【讨论】:

  • 嗨亚当,感谢您的代码。实际上,您的代码确实有所不同。您的代码中有两个线程,一个是编写器,另一个是读取器。每当作者向文件写入内容时,它都会通知读者,以便读者可以阅读新内容。但是,在我的例子中,writer 是一个 C 进程,它无法通知 Java reader 进程。
猜你喜欢
  • 2010-12-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-06-07
  • 2016-10-06
  • 2013-01-17
  • 1970-01-01
相关资源
最近更新 更多