Package org.apache.cassandra.net
Class AsyncChannelOutputPlus
java.lang.Object
java.io.OutputStream
org.apache.cassandra.io.util.DataOutputStreamPlus
org.apache.cassandra.io.util.BufferedDataOutputStreamPlus
org.apache.cassandra.net.AsyncChannelOutputPlus
- All Implemented Interfaces:
Closeable,DataOutput,Flushable,AutoCloseable,DataOutputPlus
- Direct Known Subclasses:
AsyncMessageOutputPlus,AsyncStreamingOutputPlus
A
DataOutputStreamPlus that writes ASYNCHRONOUSLY to a Netty Channel.
The close() and flush() methods synchronously wait for pending writes, and will propagate any exceptions
encountered in writing them to the wire.
The correctness of this class depends on the ChannelPromise we create against a Channel always being completed,
which appears to be a guarantee provided by Netty so long as the event loop is running.
There are two logical threads accessing the state in this class: the eventLoop of the channel, and the writer
(the writer thread may change, so long as only one utilises the class at any time).
Each thread has exclusive write access to certain state in the class, with the other thread only viewing the state,
simplifying concurrency considerations.-
Nested Class Summary
Nested Classes -
Field Summary
Fields inherited from class org.apache.cassandra.io.util.BufferedDataOutputStreamPlus
buffer -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionprotected io.netty.channel.ChannelPromisebeginFlush(long byteCount, long lowWaterMark, long highWaterMark) Create a ChannelPromise for a flush of the given size.voidclose()Flush any remaining writes, and release any buffers.abstract voiddiscard()Discard any buffered data, and the buffers that contain it.protected abstract voiddoFlush(int count) voidflush()Perform an asynchronous flush, then waits until all outstanding flushes have completedlongflushed()longprotected WritableByteChannelprotected voidparkUntilFlushed(long wakeUpWhenFlushed, long signalWhenFlushed) Utility method for waitUntilFlushed, which actually parks the current thread until the necessary number of bytes have been flushed This may only be invoked by the writer thread, never by the eventLoop.abstract longposition()Returns the current position of the underlying target like a file-pointer or the position withing a buffer.protected voidreleaseSpace(long bytesFlushed) Update our flushed count, and signal any waiters.Methods inherited from class org.apache.cassandra.io.util.BufferedDataOutputStreamPlus
allocate, order, write, write, write, write, writeBoolean, writeByte, writeBytes, writeChar, writeChars, writeDouble, writeFloat, writeInt, writeLong, writeMostSignificantBytes, writeShort, writeUTFMethods inherited from class org.apache.cassandra.io.util.DataOutputStreamPlus
retrieveTemporaryBufferMethods inherited from class java.io.OutputStream
nullOutputStreamMethods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface org.apache.cassandra.io.util.DataOutputPlus
bytesLeftInPage, hasPosition, maxBytesInPage, paddedPosition, padToPageBoundary, write, writeUnsignedVInt, writeUnsignedVInt, writeUnsignedVInt32, writeVInt, writeVInt, writeVInt32
-
Constructor Details
-
AsyncChannelOutputPlus
public AsyncChannelOutputPlus(io.netty.channel.Channel channel)
-
-
Method Details
-
beginFlush
protected io.netty.channel.ChannelPromise beginFlush(long byteCount, long lowWaterMark, long highWaterMark) throws IOException Create a ChannelPromise for a flush of the given size.This method will not return until the write is permitted by the provided watermarks and in flight bytes, and on its completion will mark the requested bytes flushed.
If this method returns normally, the ChannelPromise MUST be writtenAndFlushed, or else completed exceptionally.
- Throws:
IOException
-
parkUntilFlushed
protected void parkUntilFlushed(long wakeUpWhenFlushed, long signalWhenFlushed) Utility method for waitUntilFlushed, which actually parks the current thread until the necessary number of bytes have been flushed This may only be invoked by the writer thread, never by the eventLoop. -
releaseSpace
protected void releaseSpace(long bytesFlushed) Update our flushed count, and signal any waiters. This may only be invoked by the eventLoop, never by the writer thread. -
doFlush
- Overrides:
doFlushin classBufferedDataOutputStreamPlus- Throws:
IOException
-
position
public abstract long position()Description copied from interface:DataOutputPlusReturns the current position of the underlying target like a file-pointer or the position withing a buffer. Not every implementation may support this functionality. Whether or not this functionality is supported can be checked via theDataOutputPlus.hasPosition(). -
flushed
public long flushed() -
flushedToNetwork
public long flushedToNetwork() -
flush
Perform an asynchronous flush, then waits until all outstanding flushes have completed- Specified by:
flushin interfaceFlushable- Overrides:
flushin classBufferedDataOutputStreamPlus- Throws:
IOException- if any flush fails
-
close
Flush any remaining writes, and release any buffers. The channel is not closed, as it is assumed to be managed externally. WARNING: This method requires mutual exclusivity with all other producer methods to run safely. It should only be invoked by the owning thread, never the eventLoop; the eventLoop should propagate errors toflushFailed, which will propagate them to the producer thread no later than its final invocation toclose()orflush()(that must not be followed by any further writes).- Specified by:
closein interfaceAutoCloseable- Specified by:
closein interfaceCloseable- Overrides:
closein classBufferedDataOutputStreamPlus- Throws:
IOException
-
discard
public abstract void discard()Discard any buffered data, and the buffers that contain it. May be invoked instead ofclose()if we terminate exceptionally. -
newDefaultChannel
- Overrides:
newDefaultChannelin classDataOutputStreamPlus
-