Package org.apache.cassandra.streaming
Class StreamTransferTask
java.lang.Object
org.apache.cassandra.streaming.StreamTask
org.apache.cassandra.streaming.StreamTransferTask
StreamTransferTask sends streams for a given table
-
Field Summary
FieldsFields inherited from class org.apache.cassandra.streaming.StreamTask
session, tableId -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionvoidabort()Abort the task.voidaddTransferStream(OutgoingStream stream) voidcomplete(int sequenceNumber) Received ACK for stream atsequenceNumber.createMessageForRetry(int sequenceNumber) intlongscheduleTimeout(int sequenceNumber, long time, TimeUnit unit) Schedule timeout task to release reference for stream sent.static voidshutdownAndWait(long timeout, TimeUnit units) voidtimeout(int sequenceNumber) Received ACK for stream atsequenceNumber.Methods inherited from class org.apache.cassandra.streaming.StreamTask
getSummary
-
Field Details
-
streams
-
-
Constructor Details
-
StreamTransferTask
-
-
Method Details
-
addTransferStream
-
complete
public void complete(int sequenceNumber) Received ACK for stream atsequenceNumber.- Parameters:
sequenceNumber- sequence number of stream
-
timeout
public void timeout(int sequenceNumber) Received ACK for stream atsequenceNumber.- Parameters:
sequenceNumber- sequence number of stream
-
abort
public void abort()Description copied from class:StreamTaskAbort the task. Subclass should implement cleaning up resources.- Specified by:
abortin classStreamTask
-
getTotalNumberOfFiles
public int getTotalNumberOfFiles()- Specified by:
getTotalNumberOfFilesin classStreamTask- Returns:
- total number of files this task receives/streams.
-
getTotalSize
public long getTotalSize()- Specified by:
getTotalSizein classStreamTask- Returns:
- total bytes expected to receive
-
getFileMessages
-
createMessageForRetry
-
scheduleTimeout
Schedule timeout task to release reference for stream sent. When not receiving ACK after sending to receiver in given time, the task will release reference.- Parameters:
sequenceNumber- sequence number of stream sent.time- time to timeoutunit- unit of given time- Returns:
- scheduled future for timeout task
-
shutdownAndWait
public static void shutdownAndWait(long timeout, TimeUnit units) throws InterruptedException, TimeoutException - Throws:
InterruptedExceptionTimeoutException
-