Package org.apache.cassandra.streaming
Class StreamCoordinator
- java.lang.Object
-
- org.apache.cassandra.streaming.StreamCoordinator
-
public class StreamCoordinator extends java.lang.ObjectStreamCoordinatoris a helper class that abstracts away maintaining multiple StreamSession and ProgressInfo instances per peer. This class coordinates multiple SessionStreams per peer in both the outgoing StreamPlan context and on the inbound StreamResultFuture context.
-
-
Constructor Summary
Constructors Constructor Description StreamCoordinator(StreamOperation streamOperation, int connectionsPerHost, StreamConnectionFactory factory, boolean follower, boolean connectSequentially, java.util.UUID pendingRepair, PreviewKind previewKind)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description voidaddSessionInfo(SessionInfo session)voidconnect(StreamResultFuture future)java.util.Set<SessionInfo>getAllSessionInfo()java.util.Collection<StreamSession>getAllStreamSessions()StreamSessiongetOrCreateNextSession(InetAddressAndPort peer)StreamSessiongetOrCreateSessionById(InetAddressAndPort peer, int id)java.util.Set<InetAddressAndPort>getPeers()java.util.UUIDgetPendingRepair()StreamSessiongetSessionById(InetAddressAndPort peer, int id)booleanhasActiveSessions()booleanisFollower()voidsetConnectionFactory(StreamConnectionFactory factory)voidtransferStreams(InetAddressAndPort to, java.util.Collection<OutgoingStream> streams)voidupdateProgress(ProgressInfo info)
-
-
-
Constructor Detail
-
StreamCoordinator
public StreamCoordinator(StreamOperation streamOperation, int connectionsPerHost, StreamConnectionFactory factory, boolean follower, boolean connectSequentially, java.util.UUID pendingRepair, PreviewKind previewKind)
-
-
Method Detail
-
setConnectionFactory
public void setConnectionFactory(StreamConnectionFactory factory)
-
hasActiveSessions
public boolean hasActiveSessions()
- Returns:
- true if any stream session is active
-
getAllStreamSessions
public java.util.Collection<StreamSession> getAllStreamSessions()
-
isFollower
public boolean isFollower()
-
connect
public void connect(StreamResultFuture future)
-
getPeers
public java.util.Set<InetAddressAndPort> getPeers()
-
getOrCreateNextSession
public StreamSession getOrCreateNextSession(InetAddressAndPort peer)
-
getOrCreateSessionById
public StreamSession getOrCreateSessionById(InetAddressAndPort peer, int id)
-
getSessionById
public StreamSession getSessionById(InetAddressAndPort peer, int id)
-
updateProgress
public void updateProgress(ProgressInfo info)
-
addSessionInfo
public void addSessionInfo(SessionInfo session)
-
getAllSessionInfo
public java.util.Set<SessionInfo> getAllSessionInfo()
-
transferStreams
public void transferStreams(InetAddressAndPort to, java.util.Collection<OutgoingStream> streams)
-
getPendingRepair
public java.util.UUID getPendingRepair()
-
-