Package org.apache.cassandra.net
Class AsyncStreamingInputPlus
- java.lang.Object
-
- java.io.InputStream
-
- org.apache.cassandra.io.util.RebufferingInputStream
-
- org.apache.cassandra.net.AsyncStreamingInputPlus
-
- All Implemented Interfaces:
java.io.Closeable,java.io.DataInput,java.lang.AutoCloseable,DataInputPlus
public class AsyncStreamingInputPlus extends RebufferingInputStream
-
-
Nested Class Summary
Nested Classes Modifier and Type Class Description static interfaceAsyncStreamingInputPlus.Consumerstatic classAsyncStreamingInputPlus.InputTimeoutException-
Nested classes/interfaces inherited from interface org.apache.cassandra.io.util.DataInputPlus
DataInputPlus.DataInputStreamPlus
-
-
Field Summary
-
Fields inherited from class org.apache.cassandra.io.util.RebufferingInputStream
buffer
-
-
Constructor Summary
Constructors Constructor Description AsyncStreamingInputPlus(io.netty.channel.Channel channel)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description booleanappend(io.netty.buffer.ByteBuf buf)Append aByteBufto the end of the einternal queue.voidclose()Note: This should invoked on the consuming thread.voidconsume(AsyncStreamingInputPlus.Consumer consumer, long length)Consumes bytes in the stream until the given lengthio.netty.buffer.ByteBufAllocatorgetAllocator()booleanisEmpty()voidmaybeIssueRead()protected voidreBuffer()Implementations must implement this method to refill the buffer.voidrequestClosure()Mark this stream as closed, but do not release any of the resources.intunsafeAvailable()As long as this method is invoked on the consuming thread the returned value will be accurate.-
Methods inherited from class org.apache.cassandra.io.util.RebufferingInputStream
read, read, readBoolean, readByte, readChar, readDouble, readFloat, readFully, readFully, readFully, readInt, readLine, readLong, readPrimitiveSlowly, readShort, readUnsignedByte, readUnsignedShort, readUnsignedVInt, readUTF, readVInt, skipBytes
-
Methods inherited from class java.io.InputStream
available, mark, markSupported, nullInputStream, read, readAllBytes, readNBytes, readNBytes, reset, skip, transferTo
-
Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
-
Methods inherited from interface org.apache.cassandra.io.util.DataInputPlus
skipBytesFully
-
-
-
-
Method Detail
-
append
public boolean append(io.netty.buffer.ByteBuf buf) throws java.lang.IllegalStateExceptionAppend aByteBufto the end of the einternal queue. Note: it's expected this method is invoked on the netty event loop.- Throws:
java.lang.IllegalStateException
-
reBuffer
protected void reBuffer() throws java.io.EOFException, AsyncStreamingInputPlus.InputTimeoutExceptionImplementations must implement this method to refill the buffer. They can expect the buffer to be empty when this method is invoked. Release open buffers and poll thequeuefor more data.This is best, and more or less expected, to be invoked on a consuming thread (not the event loop) becasue if we block on the queue we can't fill it on the event loop (as that's where the buffers are coming from).
- Specified by:
reBufferin classRebufferingInputStream- Throws:
java.io.EOFException- when no further reading from this instance should occur. Implies this instance is closed.AsyncStreamingInputPlus.InputTimeoutException- when no new buffers arrive for reading before therebufferTimeoutNanoselapses while blocking. It's then not safe to reuse this instance again.
-
consume
public void consume(AsyncStreamingInputPlus.Consumer consumer, long length) throws java.io.IOException
Consumes bytes in the stream until the given length- Throws:
java.io.IOException
-
unsafeAvailable
public int unsafeAvailable()
As long as this method is invoked on the consuming thread the returned value will be accurate.
-
maybeIssueRead
public void maybeIssueRead()
-
isEmpty
public boolean isEmpty()
-
close
public void close()
Note: This should invoked on the consuming thread.- Specified by:
closein interfacejava.lang.AutoCloseable- Specified by:
closein interfacejava.io.Closeable- Overrides:
closein classjava.io.InputStream
-
requestClosure
public void requestClosure()
Mark this stream as closed, but do not release any of the resources. Note: this is best to be called from the producer thread.
-
getAllocator
public io.netty.buffer.ByteBufAllocator getAllocator()
-
-