From 6ad995e8fa7703e082ea9ce67dc4c1ed0b1fd18a Mon Sep 17 00:00:00 2001 From: Jason Brown Date: Thu, 30 Jan 2014 09:51:42 -0800 Subject: [PATCH] sstables from stalled repair sessions become live after a reboot and can resurrect deleted data patch by jasobrown, reviewed by yukim for CASSANDRA-6503 --- CHANGES.txt | 1 + .../cassandra/streaming/IncomingStreamReader.java | 8 ++++---- .../apache/cassandra/streaming/StreamInSession.java | 13 +++++++++---- 3 files changed, 14 insertions(+), 8 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 110bf505e8..d85d3a4b69 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -21,6 +21,7 @@ * Fix preparing with batch and delete from collection (CASSANDRA-6607) * Fix ABSC reverse iterator's remove() method (CASSANDRA-6629) * Handle host ID conflicts properly (CASSANDRA-6615) + * sstables from stalled repair sessions can resurrect deleted data (CASSANDRA-6503) 1.2.13 diff --git a/src/java/org/apache/cassandra/streaming/IncomingStreamReader.java b/src/java/org/apache/cassandra/streaming/IncomingStreamReader.java index 0b058fc72b..940f8de5cf 100644 --- a/src/java/org/apache/cassandra/streaming/IncomingStreamReader.java +++ b/src/java/org/apache/cassandra/streaming/IncomingStreamReader.java @@ -119,8 +119,8 @@ public class IncomingStreamReader DataInput dis = new DataInputStream(underliningStream); try { - SSTableReader reader = streamIn(dis, localFile, remoteFile); - session.finished(remoteFile, reader); + SSTableWriter writer = streamIn(dis, localFile, remoteFile); + session.finished(remoteFile, writer); } catch (IOException ex) { @@ -141,7 +141,7 @@ public class IncomingStreamReader /** * @throws IOException if reading the remote sstable fails. Will throw an RTE if local write fails. */ - private SSTableReader streamIn(DataInput input, PendingFile localFile, PendingFile remoteFile) throws IOException + private SSTableWriter streamIn(DataInput input, PendingFile localFile, PendingFile remoteFile) throws IOException { ColumnFamilyStore cfs = Table.open(localFile.desc.ksname).getColumnFamilyStore(localFile.desc.cfname); DecoratedKey key; @@ -197,7 +197,7 @@ public class IncomingStreamReader } StreamingMetrics.totalIncomingBytes.inc(totalBytesRead); metrics.incomingBytes.inc(totalBytesRead); - return writer.closeAndOpenReader(); + return writer; } catch (Throwable e) { diff --git a/src/java/org/apache/cassandra/streaming/StreamInSession.java b/src/java/org/apache/cassandra/streaming/StreamInSession.java index 96c31da0f4..e83a5b6405 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.net.Socket; import java.util.*; import java.util.concurrent.ConcurrentMap; +import org.apache.cassandra.io.sstable.SSTableWriter; import org.cliffc.high_scale_lib.NonBlockingHashMap; import org.cliffc.high_scale_lib.NonBlockingHashSet; import org.slf4j.Logger; @@ -47,7 +48,7 @@ public class StreamInSession extends AbstractStreamSession private static final ConcurrentMap sessions = new NonBlockingHashMap(); private final Set files = new NonBlockingHashSet(); - private final List readers = new ArrayList(); + private final List writers = new ArrayList(); private PendingFile current; private Socket socket; private volatile int retries; @@ -106,13 +107,13 @@ public class StreamInSession extends AbstractStreamSession } } - public void finished(PendingFile remoteFile, SSTableReader reader) throws IOException + public void finished(PendingFile remoteFile, SSTableWriter writer) throws IOException { if (logger.isDebugEnabled()) logger.debug("Finished {} (from {}). Sending ack to {}", new Object[] {remoteFile, getHost(), this}); - assert reader != null; - readers.add(reader); + assert writer != null; + writers.add(writer); files.remove(remoteFile); if (remoteFile.equals(current)) current = null; @@ -163,6 +164,10 @@ public class StreamInSession extends AbstractStreamSession HashMap > cfstores = new HashMap>(); try { + List readers = new ArrayList(); + for(SSTableWriter writer : writers) + readers.add(writer.closeAndOpenReader()); + for (SSTableReader sstable : readers) { assert sstable.getTableName().equals(table);