Class StreamingMultiplexedChannel

java.lang.Object
org.apache.cassandra.streaming.async.StreamingMultiplexedChannel

public class StreamingMultiplexedChannel extends Object
Responsible for sending StreamMessages to a given peer. We manage an array of netty Channels for sending OutgoingStreamMessage instances; all other StreamMessage types are sent via a special control channel. The reason for this is to treat those messages carefully and not let them get stuck behind a stream transfer. One of the challenges when sending streams is we might need to delay shipping the stream if: - we've exceeded our network I/O use due to rate limiting (at the cassandra level) - the receiver isn't keeping up, which causes the local TCP socket buffer to not empty, which causes epoll writes to not move any bytes to the socket, which causes buffers to stick around in user-land (a/k/a cassandra) memory. When those conditions occur, it's easy enough to reschedule processing the stream once the resources pick up (we acquire the permits from the rate limiter, or the socket drains). However, we need to ensure that no other messages are submitted to the same channel while the current stream is still being processed.
  • Constructor Details

  • Method Details

    • peer

      public InetAddressAndPort peer()
    • connectedTo

      public InetSocketAddress connectedTo()
    • sendControlMessage

      public io.netty.util.concurrent.Future<?> sendControlMessage(StreamMessage message)
    • sendMessage

      public io.netty.util.concurrent.Future<?> sendMessage(StreamingChannel channel, StreamMessage message)
    • setClosed

      public void setClosed()
      For testing purposes only.
    • connected

      public boolean connected()
    • close

      public void close()
    • unsafeCloseControlChannel

      public void unsafeCloseControlChannel()