public class ReadAheadInputStream
extends java.io.FilterInputStream
InputStream
to asynchronously read ahead from an underlying input stream when a specified amount of data has been read from the current
buffer. It does so by maintaining two buffers: an active buffer and a read ahead buffer. The active buffer contains data which should be returned when a
read() call is issued. The read ahead buffer is used to asynchronously read from the underlying input stream. When the current active buffer is exhausted, we
flip the two buffers so that we can start reading from the read ahead buffer without being blocked by disk I/O.
To build an instance, see ReadAheadInputStream.Builder
.
This class was ported and adapted from Apache Spark commit 933dc6cb7b3de1d8ccaf73d124d6eb95b947ed19.
Modifier and Type | Class and Description |
---|---|
static class |
ReadAheadInputStream.Builder
Builds a new
ReadAheadInputStream instance. |
Modifier and Type | Field and Description |
---|---|
private java.nio.ByteBuffer |
activeBuffer |
private java.util.concurrent.locks.Condition |
asyncReadComplete |
private static java.lang.ThreadLocal<byte[]> |
BYTE_ARRAY_1 |
private boolean |
endOfStream |
private java.util.concurrent.ExecutorService |
executorService |
private boolean |
isClosed |
private boolean |
isReading |
private boolean |
isUnderlyingInputStreamBeingClosed |
private java.util.concurrent.atomic.AtomicBoolean |
isWaiting |
private boolean |
readAborted |
private java.nio.ByteBuffer |
readAheadBuffer |
private java.lang.Throwable |
readException |
private boolean |
readInProgress |
private boolean |
shutdownExecutorService |
private java.util.concurrent.locks.ReentrantLock |
stateChangeLock |
Modifier | Constructor and Description |
---|---|
|
ReadAheadInputStream(java.io.InputStream inputStream,
int bufferSizeInBytes)
Deprecated.
|
|
ReadAheadInputStream(java.io.InputStream inputStream,
int bufferSizeInBytes,
java.util.concurrent.ExecutorService executorService)
Deprecated.
|
private |
ReadAheadInputStream(java.io.InputStream inputStream,
int bufferSizeInBytes,
java.util.concurrent.ExecutorService executorService,
boolean shutdownExecutorService)
Constructs an instance with the specified buffer size and read-ahead threshold
|
Modifier and Type | Method and Description |
---|---|
int |
available() |
static ReadAheadInputStream.Builder |
builder()
Constructs a new
ReadAheadInputStream.Builder . |
private void |
checkReadException() |
void |
close() |
private void |
closeUnderlyingInputStreamIfNecessary() |
private boolean |
isEndOfStream() |
private static java.lang.Thread |
newDaemonThread(java.lang.Runnable r)
Constructs a new daemon thread.
|
private static java.util.concurrent.ExecutorService |
newExecutorService()
Constructs a new daemon executor service.
|
int |
read() |
int |
read(byte[] b,
int offset,
int len) |
private void |
readAsync()
Read data from underlyingInputStream to readAheadBuffer asynchronously.
|
private void |
signalAsyncReadComplete() |
long |
skip(long n) |
private long |
skipInternal(long n)
Internal skip function which should be called only from skip().
|
private void |
swapBuffers()
Flips the active and read ahead buffer
|
private void |
waitForAsyncReadComplete() |
private static final java.lang.ThreadLocal<byte[]> BYTE_ARRAY_1
private final java.util.concurrent.locks.ReentrantLock stateChangeLock
private java.nio.ByteBuffer activeBuffer
private java.nio.ByteBuffer readAheadBuffer
private boolean endOfStream
private boolean readInProgress
private boolean readAborted
private java.lang.Throwable readException
private boolean isClosed
private boolean isUnderlyingInputStreamBeingClosed
private boolean isReading
private final java.util.concurrent.atomic.AtomicBoolean isWaiting
private final java.util.concurrent.ExecutorService executorService
private final boolean shutdownExecutorService
private final java.util.concurrent.locks.Condition asyncReadComplete
@Deprecated public ReadAheadInputStream(java.io.InputStream inputStream, int bufferSizeInBytes)
builder()
, ReadAheadInputStream.Builder
, and ReadAheadInputStream.Builder.get()
inputStream
- The underlying input stream.bufferSizeInBytes
- The buffer size.@Deprecated public ReadAheadInputStream(java.io.InputStream inputStream, int bufferSizeInBytes, java.util.concurrent.ExecutorService executorService)
builder()
, ReadAheadInputStream.Builder
, and ReadAheadInputStream.Builder.get()
inputStream
- The underlying input stream.bufferSizeInBytes
- The buffer size.executorService
- An executor service for the read-ahead thread.private ReadAheadInputStream(java.io.InputStream inputStream, int bufferSizeInBytes, java.util.concurrent.ExecutorService executorService, boolean shutdownExecutorService)
inputStream
- The underlying input stream.bufferSizeInBytes
- The buffer size.executorService
- An executor service for the read-ahead thread.shutdownExecutorService
- Whether or not to shut down the given ExecutorService on close.public static ReadAheadInputStream.Builder builder()
ReadAheadInputStream.Builder
.ReadAheadInputStream.Builder
.private static java.lang.Thread newDaemonThread(java.lang.Runnable r)
r
- the thread's runnable.private static java.util.concurrent.ExecutorService newExecutorService()
public int available() throws java.io.IOException
available
in class java.io.FilterInputStream
java.io.IOException
private void checkReadException() throws java.io.IOException
java.io.IOException
public void close() throws java.io.IOException
close
in interface java.io.Closeable
close
in interface java.lang.AutoCloseable
close
in class java.io.FilterInputStream
java.io.IOException
private void closeUnderlyingInputStreamIfNecessary()
private boolean isEndOfStream()
public int read() throws java.io.IOException
read
in class java.io.FilterInputStream
java.io.IOException
public int read(byte[] b, int offset, int len) throws java.io.IOException
read
in class java.io.FilterInputStream
java.io.IOException
private void readAsync() throws java.io.IOException
java.io.IOException
- if an I/O error occurs.private void signalAsyncReadComplete()
public long skip(long n) throws java.io.IOException
skip
in class java.io.FilterInputStream
java.io.IOException
private long skipInternal(long n) throws java.io.IOException
n
- the number of bytes to be skipped.java.io.IOException
- if an I/O error occurs.private void swapBuffers()
private void waitForAsyncReadComplete() throws java.io.IOException
java.io.IOException