Package org.apache.cassandra.net
Class AsyncStreamingInputPlus
java.lang.Object
java.io.InputStream
org.apache.cassandra.io.util.DataInputPlus.DataInputStreamPlus
org.apache.cassandra.io.util.RebufferingInputStream
org.apache.cassandra.net.AsyncStreamingInputPlus
- All Implemented Interfaces:
Closeable,DataInput,AutoCloseable,DataInputPlus,StreamingDataInputPlus
public class AsyncStreamingInputPlus
extends RebufferingInputStream
implements StreamingDataInputPlus
-
Nested Class Summary
Nested ClassesNested 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 -
Method Summary
Modifier and TypeMethodDescriptionbooleanappend(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.ByteBufAllocatorbooleanisEmpty()protected voidreBuffer()Implementations must implement this method to refill the buffer.voidMark this stream as closed, but do not release any of the resources.intAs 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, readUnsignedVInt32, readUTF, readVInt, readVInt32, skipBytesMethods inherited from class java.io.InputStream
available, mark, markSupported, nullInputStream, read, readAllBytes, readNBytes, readNBytes, reset, skip, skipNBytes, transferToMethods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface java.io.DataInput
readBoolean, readByte, readChar, readDouble, readFloat, readFully, readFully, readInt, readLine, readLong, readShort, readUnsignedByte, readUnsignedShort, readUTFMethods inherited from interface org.apache.cassandra.io.util.DataInputPlus
readUnsignedVInt, readUnsignedVInt32, readVInt, readVInt32, skipBytes, skipBytesFully
-
Constructor Details
-
AsyncStreamingInputPlus
public AsyncStreamingInputPlus(io.netty.channel.Channel channel)
-
-
Method Details
-
append
Append aByteBufto the end of the einternal queue. Note: it's expected this method is invoked on the netty event loop.- Throws:
IllegalStateException
-
reBuffer
Implementations 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 invoked on a consuming thread (not the event loop) because 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:
ClosedChannelException- when no further reading from this instance should occur. Implies this instance is closed.
-
consume
Consumes bytes in the stream until the given length- Throws:
IOException
-
unsafeAvailable
public int unsafeAvailable()As long as this method is invoked on the consuming thread the returned value will be accurate. -
isEmpty
public boolean isEmpty() -
close
public void close()Note: This should invoked on the consuming thread.- Specified by:
closein interfaceAutoCloseable- Specified by:
closein interfaceCloseable- Specified by:
closein interfaceStreamingDataInputPlus- Overrides:
closein classInputStream
-
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()
-