From bc15e804530e87b1c5cc8c96fdf939dcaaa8efd9 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Tue, 14 Jun 2011 15:16:04 +0000 Subject: [PATCH] backport #2766 from 0.8 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1135638 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../org/apache/cassandra/streaming/StreamInSession.java | 6 ++++-- 2 files changed, 5 insertions(+), 2 deletions(-) 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)