diff --git a/CHANGES.txt b/CHANGES.txt index 994152eef6..fae47b6b09 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -18,6 +18,7 @@ * fix nodetool ring use with Ec2Snitch (CASSANDRA-2733) * fix removing columns and subcolumns that are supressed by a row or supercolumn tombstone during replica resolution (CASSANDRA-2590) + * use threadsafe collections for StreamInSession (CASSANDRA-2766) 0.7.6 diff --git a/src/java/org/apache/cassandra/streaming/StreamInSession.java b/src/java/org/apache/cassandra/streaming/StreamInSession.java index 43ee05ba31..63244aed66 100644 --- a/src/java/org/apache/cassandra/streaming/StreamInSession.java +++ b/src/java/org/apache/cassandra/streaming/StreamInSession.java @@ -24,6 +24,7 @@ import java.util.*; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; +import java.util.concurrent.LinkedBlockingQueue; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -35,6 +36,7 @@ import org.apache.cassandra.io.sstable.SSTableReader; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.utils.Pair; import org.cliffc.high_scale_lib.NonBlockingHashMap; +import org.cliffc.high_scale_lib.NonBlockingHashSet; /** each context gets its own StreamInSession. So there may be >1 Session per host */ public class StreamInSession @@ -43,11 +45,11 @@ public class StreamInSession private static ConcurrentMap, StreamInSession> sessions = new NonBlockingHashMap, StreamInSession>(); - private final List files = new ArrayList(); + private final Set files = new NonBlockingHashSet(); private final Pair context; private final Runnable callback; private String table; - private final List> buildFutures = new ArrayList>(); + private final Collection> buildFutures = new LinkedBlockingQueue>(); private PendingFile current; private StreamInSession(Pair context, Runnable callback)