From aadfa6a334135e1f9f7b8eff827941ab6e588d2f Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 18 Jun 2010 10:29:54 +0000 Subject: [PATCH] Stream sstables without anticompaction patch by Stu Hood; reviewed by jbellis for CASSANDRA-579 git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@955923 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + NEWS.txt | 2 + src/java/org/apache/cassandra/db/Table.java | 29 +--- .../apache/cassandra/dht/AbstractBounds.java | 24 +++- .../cassandra/io/sstable/SSTableReader.java | 134 ++++++++++-------- .../cassandra/io/sstable/SSTableScanner.java | 2 +- .../cassandra/io/sstable/SSTableWriter.java | 2 +- .../apache/cassandra/net/FileStreamTask.java | 31 ++-- .../cassandra/net/MessagingService.java | 8 +- .../cassandra/service/AntiEntropyService.java | 6 +- .../cassandra/streaming/FileStatus.java | 13 +- .../streaming/FileStatusHandler.java | 1 + .../streaming/IncomingStreamReader.java | 33 +++-- .../cassandra/streaming/PendingFile.java | 60 +++++--- .../streaming/StreamFinishedVerbHandler.java | 2 +- .../streaming/StreamInitiateMessage.java | 12 +- .../streaming/StreamInitiateVerbHandler.java | 4 +- .../apache/cassandra/streaming/StreamOut.java | 26 ++-- .../cassandra/streaming/StreamOutManager.java | 32 ++--- .../cassandra/streaming/StreamingService.java | 8 +- .../cassandra/db/ColumnFamilyStoreTest.java | 2 +- .../org/apache/cassandra/db/TableTest.java | 2 +- .../apache/cassandra/io/StreamingTest.java | 22 ++- .../io/sstable/LegacySSTableTest.java | 2 +- .../io/sstable/SSTableReaderTest.java | 54 ++++++- .../cassandra/io/sstable/SSTableTest.java | 4 +- .../cassandra/streaming/BootstrapTest.java | 5 +- 27 files changed, 307 insertions(+), 214 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 7a06828db4..4fdcdb187b 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -29,6 +29,7 @@ dev segments that always contain entire entries/rows (CASSANDRA-1117) * avoid reading large rows into memory during compaction (CASSANDRA-16) * added hadoop OutputFormat (CASSANDRA-1101) + * efficient Streaming (no more anticompaction) (CASSANDRA-579) 0.6.3 diff --git a/NEWS.txt b/NEWS.txt index 2c8168bafd..5c92a5a097 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -15,6 +15,8 @@ Features `cassandra.yaml.` - row size limit increased from 2GB to 2 billion columns - Hadoop OutputFormat support + - Streaming data for repair or node movement no longer requires + anticompaction step first Configuraton ------------ diff --git a/src/java/org/apache/cassandra/db/Table.java b/src/java/org/apache/cassandra/db/Table.java index 724cd177ef..96977fa907 100644 --- a/src/java/org/apache/cassandra/db/Table.java +++ b/src/java/org/apache/cassandra/db/Table.java @@ -191,40 +191,21 @@ public class Table } } } - - /* - * This method is invoked only during a bootstrap process. We basically - * do a complete compaction since we can figure out based on the ranges - * whether the files need to be split. - */ - public List forceAntiCompaction(Collection ranges, InetAddress target) - { - List allResults = new ArrayList(); - for (ColumnFamilyStore cfStore : columnFamilyStores.values()) - { - try - { - allResults.addAll(CompactionManager.instance.submitAnticompaction(cfStore, ranges, target).get()); - } - catch (Exception e) - { - throw new RuntimeException(e); - } - } - return allResults; - } /* * This method is an ADMIN operation to force compaction * of all SSTables on disk. - */ + */ public void forceCompaction() { for (ColumnFamilyStore cfStore : columnFamilyStores.values()) CompactionManager.instance.submitMajor(cfStore); } - List getAllSSTablesOnDisk() + /** + * @return A list of open SSTableReaders (TODO: ensure that the caller doesn't modify these). + */ + public List getAllSSTables() { List list = new ArrayList(); for (ColumnFamilyStore cfStore : columnFamilyStores.values()) diff --git a/src/java/org/apache/cassandra/dht/AbstractBounds.java b/src/java/org/apache/cassandra/dht/AbstractBounds.java index 85ab4577aa..807a1979fc 100644 --- a/src/java/org/apache/cassandra/dht/AbstractBounds.java +++ b/src/java/org/apache/cassandra/dht/AbstractBounds.java @@ -25,8 +25,7 @@ import java.io.DataInput; import java.io.DataOutput; import java.io.IOException; import java.io.Serializable; -import java.util.List; -import java.util.Set; +import java.util.*; import org.apache.cassandra.io.ICompactSerializer2; @@ -72,6 +71,27 @@ public abstract class AbstractBounds implements Serializable public abstract List unwrap(); + /** + * @return A copy of the given list of non-intersecting bounds with all bounds unwrapped, sorted by bound.left. + */ + public static List normalize(Collection bounds) + { + // unwrap all + List output = new ArrayList(); + for (AbstractBounds bound : bounds) + output.addAll(bound.unwrap()); + + // sort by left + Collections.sort(output, new Comparator() + { + public int compare(AbstractBounds b1, AbstractBounds b2) + { + return b1.left.compareTo(b2.left); + } + }); + return output; + } + private static class AbstractBoundsSerializer implements ICompactSerializer2 { public void serialize(AbstractBounds range, DataOutput out) throws IOException diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java index 942fbb49b9..c96174100e 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java @@ -42,9 +42,12 @@ import org.apache.cassandra.db.*; import org.apache.cassandra.db.clock.AbstractReconciler; import org.apache.cassandra.db.filter.QueryFilter; import org.apache.cassandra.db.marshal.AbstractType; +import org.apache.cassandra.dht.AbstractBounds; +import org.apache.cassandra.dht.Range; import org.apache.cassandra.io.ICompactSerializer2; import org.apache.cassandra.io.util.FileDataInput; import org.apache.cassandra.utils.BloomFilter; +import org.apache.cassandra.utils.Pair; /** * SSTableReaders are open()ed by Table.onStart; after that they are created by SSTableWriter.renameAndOpen. @@ -344,15 +347,40 @@ public class SSTableReader extends SSTable implements Comparable } /** - * Returns the position in the data file to find the given key, or -1 if the - * key is not present. - * FIXME: should not be public: use Scanner. + * Determine the minimal set of sections that can be extracted from this SSTable to cover the given ranges. + * @return A sorted list of (offset,end) pairs that cover the given ranges in the datafile for this SSTable. */ - @Deprecated - public long getPosition(DecoratedKey decoratedKey) + public List> getPositionsForRanges(Collection ranges) + { + // use the index to determine a minimal section for each range + List> positions = new ArrayList>(); + for (AbstractBounds range : AbstractBounds.normalize(ranges)) + { + long left = getPosition(new DecoratedKey(range.left, null), Operator.GT); + if (left == -1) + // left is past the end of the file + continue; + long right = getPosition(new DecoratedKey(range.right, null), Operator.GT); + if (right == -1 || Range.isWrapAround(range.left, range.right)) + // right is past the end of the file, or it wraps + right = length(); + if (left == right) + // empty range + continue; + positions.add(new Pair(Long.valueOf(left), Long.valueOf(right))); + } + return positions; + } + + /** + * @param decoratedKey The key to apply as the rhs to the given Operator. + * @param op The Operator defining matching keys: the nearest key to the target matching the operator wins. + * @return The position in the data file to find the key, or -1 if the key is not present + */ + public long getPosition(DecoratedKey decoratedKey, Operator op) { // first, check bloom filter - if (!bf.isPresent(partitioner.convertToDiskFormat(decoratedKey))) + if (op == Operator.EQ && !bf.isPresent(partitioner.convertToDiskFormat(decoratedKey))) return -1; // next, the key cache @@ -369,30 +397,32 @@ public class SSTableReader extends SSTable implements Comparable // next, see if the sampled index says it's impossible for the key to be present IndexSummary.KeyPosition sampledPosition = getIndexScanPosition(decoratedKey); if (sampledPosition == null) - return -1; + // we matched the -1th position: if the operator might match forward, return the 0th position + return op.apply(1) >= 0 ? 0 : -1; // scan the on-disk index, starting at the nearest sampled position - int i = 0; Iterator segments = ifile.iterator(sampledPosition.indexPosition, INDEX_FILE_BUFFER_BYTES); while (segments.hasNext()) { FileDataInput input = segments.next(); try { - while (!input.isEOF() && i++ < IndexSummary.INDEX_INTERVAL) + while (!input.isEOF()) { // read key & data position from index entry DecoratedKey indexDecoratedKey = partitioner.convertFromDiskFormat(FBUtilities.readShortByteArray(input)); long dataPosition = input.readLong(); - int v = indexDecoratedKey.compareTo(decoratedKey); + int comparison = indexDecoratedKey.compareTo(decoratedKey); + int v = op.apply(comparison); if (v == 0) { - if (keyCache != null && keyCache.getCapacity() > 0) + if (comparison == 0 && keyCache != null && keyCache.getCapacity() > 0) + // store exact match for the key keyCache.put(unifiedKey, Long.valueOf(dataPosition)); return dataPosition; } - if (v > 0) + if (v < 0) return -1; } } @@ -415,53 +445,6 @@ public class SSTableReader extends SSTable implements Comparable return -1; } - /** - * Like getPosition, but if key is not found will return the location of the - * first key _greater_ than the desired one, or -1 if no such key exists. - * FIXME: should not be public: use Scanner. - */ - @Deprecated - public long getNearestPosition(DecoratedKey decoratedKey) - { - IndexSummary.KeyPosition sampledPosition = getIndexScanPosition(decoratedKey); - if (sampledPosition == null) - return 0; - - // scan the on-disk index, starting at the nearest sampled position - Iterator segiter = ifile.iterator(sampledPosition.indexPosition, INDEX_FILE_BUFFER_BYTES); - while (segiter.hasNext()) - { - FileDataInput input = segiter.next(); - try - { - while (!input.isEOF()) - { - DecoratedKey indexDecoratedKey = partitioner.convertFromDiskFormat(FBUtilities.readShortByteArray(input)); - long position = input.readLong(); - int v = indexDecoratedKey.compareTo(decoratedKey); - if (v >= 0) - return position; - } - } - catch (IOException e) - { - throw new IOError(e); - } - finally - { - try - { - input.close(); - } - catch (IOException e) - { - logger.error("error closing file", e); - } - } - } - return -1; - } - /** * @return The length in bytes of the data file for this SSTable. */ @@ -507,7 +490,7 @@ public class SSTableReader extends SSTable implements Comparable public FileDataInput getFileDataInput(DecoratedKey decoratedKey, int bufferSize) { - long position = getPosition(decoratedKey); + long position = getPosition(decoratedKey, Operator.EQ); if (position < 0) return null; @@ -557,4 +540,35 @@ public class SSTableReader extends SSTable implements Comparable return in.readInt(); return in.readLong(); } + + /** + * TODO: Move someplace reusable + */ + public abstract static class Operator + { + public static final Operator EQ = new Equals(); + public static final Operator GE = new GreaterThanOrEqualTo(); + public static final Operator GT = new GreaterThan(); + + /** + * @param comparison The result of a call to compare/compareTo, with the desired field on the rhs. + * @return less than 0 if the operator cannot match forward, 0 if it matches, greater than 0 if it might match forward. + */ + public abstract int apply(int comparison); + + final static class Equals extends Operator + { + public int apply(int comparison) { return -comparison; } + } + + final static class GreaterThanOrEqualTo extends Operator + { + public int apply(int comparison) { return comparison >= 0 ? 0 : -comparison; } + } + + final static class GreaterThan extends Operator + { + public int apply(int comparison) { return comparison > 0 ? 0 : 1; } + } + } } diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java b/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java index aec52227b8..b9a49e61c2 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableScanner.java @@ -90,7 +90,7 @@ public class SSTableScanner implements Iterator, Closeable { try { - long position = sstable.getNearestPosition(seekKey); + long position = sstable.getPosition(seekKey, SSTableReader.Operator.GE); if (position < 0) { exhausted = true; diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java index 6d5187e676..9c3fd54137 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableWriter.java @@ -39,13 +39,13 @@ package org.apache.cassandra.io.sstable; import java.io.*; -import org.apache.cassandra.io.AbstractCompactedRow; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.DecoratedKey; import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.io.AbstractCompactedRow; import org.apache.cassandra.io.util.BufferedRandomAccessFile; import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.io.util.SegmentedFile; diff --git a/src/java/org/apache/cassandra/net/FileStreamTask.java b/src/java/org/apache/cassandra/net/FileStreamTask.java index 48e191d766..bbd05b75bb 100644 --- a/src/java/org/apache/cassandra/net/FileStreamTask.java +++ b/src/java/org/apache/cassandra/net/FileStreamTask.java @@ -25,12 +25,13 @@ import java.nio.ByteBuffer; import java.nio.channels.FileChannel; import java.nio.channels.SocketChannel; -import org.apache.cassandra.streaming.StreamOutManager; +import org.apache.cassandra.streaming.PendingFile; import org.apache.cassandra.utils.FBUtilities; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.WrappedRunnable; public class FileStreamTask extends WrappedRunnable @@ -41,21 +42,12 @@ public class FileStreamTask extends WrappedRunnable // around 10 minutes at the default rpctimeout public static final int MAX_CONNECT_ATTEMPTS = 8; - private final String file; - private final long startPosition; - private final long endPosition; + private final PendingFile file; private final InetAddress to; - FileStreamTask(String file, InetAddress to) - { - this(file, 0, new File(file).length(), to); - } - - private FileStreamTask(String file, long startPosition, long endPosition, InetAddress to) + FileStreamTask(PendingFile file, InetAddress to) { this.file = file; - this.startPosition = startPosition; - this.endPosition = endPosition; this.to = to; } @@ -87,8 +79,7 @@ public class FileStreamTask extends WrappedRunnable private void stream(SocketChannel channel) throws IOException { - long start = startPosition; - RandomAccessFile raf = new RandomAccessFile(new File(file), "r"); + RandomAccessFile raf = new RandomAccessFile(new File(file.getFilename()), "r"); try { FileChannel fc = raf.getChannel(); @@ -96,14 +87,16 @@ public class FileStreamTask extends WrappedRunnable ByteBuffer buffer = MessagingService.constructStreamHeader(false); channel.write(buffer); assert buffer.remaining() == 0; - - while (start < endPosition) + + // stream sections of the file as returned by PendingFile.currentSection + Pair section; + while ((section = file.currentSection()) != null) { - long bytesTransferred = fc.transferTo(start, CHUNK_SIZE, channel); + long length = Math.min(CHUNK_SIZE, section.right - section.left); + long bytesTransferred = fc.transferTo(section.left, length, channel); if (logger.isDebugEnabled()) logger.debug("Bytes transferred " + bytesTransferred); - start += bytesTransferred; - StreamOutManager.get(to).update(file, start); + file.update(section.left + bytesTransferred); } } finally diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index dc2c59104d..474d7c2b8f 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -27,6 +27,7 @@ import org.apache.cassandra.io.util.DataOutputBuffer; import org.apache.cassandra.net.io.SerializerType; import org.apache.cassandra.net.sink.SinkManager; import org.apache.cassandra.service.StorageService; +import org.apache.cassandra.streaming.PendingFile; import org.apache.cassandra.utils.ExpiringMap; import org.apache.cassandra.utils.GuidGenerator; import org.apache.cassandra.utils.SimpleCondition; @@ -325,15 +326,14 @@ public class MessagingService implements IFailureDetectionEventListener /** * Stream a file from source to destination. This is highly optimized * to not hold any of the contents of the file in memory. - * @param file name of file to stream. + * @param file file to stream. * @param to endpoint to which we need to stream the file. */ - public void stream(String file, InetAddress to) + public void stream(PendingFile file, InetAddress to) { /* Streaming asynchronously on streamExector_ threads. */ - Runnable streamingTask = new FileStreamTask(file, to); - streamExecutor_.execute(streamingTask); + streamExecutor_.execute(new FileStreamTask(file, to)); } /** blocks until the processing pools are empty and done. */ diff --git a/src/java/org/apache/cassandra/service/AntiEntropyService.java b/src/java/org/apache/cassandra/service/AntiEntropyService.java index a5befe394e..9e94389e7c 100644 --- a/src/java/org/apache/cassandra/service/AntiEntropyService.java +++ b/src/java/org/apache/cassandra/service/AntiEntropyService.java @@ -660,13 +660,13 @@ public class AntiEntropyService ColumnFamilyStore cfstore = Table.open(cf.left).getColumnFamilyStore(cf.right); try { - List ranges = new ArrayList(differences); - final List sstables = CompactionManager.instance.submitAnticompaction(cfstore, ranges, remote).get(); + final List ranges = new ArrayList(differences); + final Collection sstables = cfstore.getSSTables(); Future f = StageManager.getStage(StageManager.STREAM_STAGE).submit(new WrappedRunnable() { protected void runMayThrow() throws Exception { - StreamOut.transferSSTables(remote, sstables, cf.left); + StreamOut.transferSSTables(remote, cf.left, sstables, ranges); StreamOutManager.remove(remote); } }); diff --git a/src/java/org/apache/cassandra/streaming/FileStatus.java b/src/java/org/apache/cassandra/streaming/FileStatus.java index 214cae888b..bfedd3c551 100644 --- a/src/java/org/apache/cassandra/streaming/FileStatus.java +++ b/src/java/org/apache/cassandra/streaming/FileStatus.java @@ -54,16 +54,14 @@ class FileStatus } private final String file_; - private final long expectedBytes_; private Action action_; /** * Create a FileStatus with the default Action: STREAM. */ - public FileStatus(String file, long expectedBytes) + public FileStatus(String file) { file_ = file; - expectedBytes_ = expectedBytes; action_ = Action.STREAM; } @@ -72,11 +70,6 @@ class FileStatus return file_; } - public long getExpectedBytes() - { - return expectedBytes_; - } - public void setAction(Action action) { action_ = action; @@ -100,15 +93,13 @@ class FileStatus public void serialize(FileStatus streamStatus, DataOutputStream dos) throws IOException { dos.writeUTF(streamStatus.getFile()); - dos.writeLong(streamStatus.getExpectedBytes()); dos.writeInt(streamStatus.getAction().ordinal()); } public FileStatus deserialize(DataInputStream dis) throws IOException { String targetFile = dis.readUTF(); - long expectedBytes = dis.readLong(); - FileStatus streamStatus = new FileStatus(targetFile, expectedBytes); + FileStatus streamStatus = new FileStatus(targetFile); int ordinal = dis.readInt(); if (ordinal == Action.DELETE.ordinal()) diff --git a/src/java/org/apache/cassandra/streaming/FileStatusHandler.java b/src/java/org/apache/cassandra/streaming/FileStatusHandler.java index c31d6edc91..85478a4f46 100644 --- a/src/java/org/apache/cassandra/streaming/FileStatusHandler.java +++ b/src/java/org/apache/cassandra/streaming/FileStatusHandler.java @@ -64,6 +64,7 @@ class FileStatusHandler } catch (IOException e) { + logger.error("Failed adding " + pendingFile, e); throw new RuntimeException("Not able to add streamed file " + pendingFile.getFilename(), e); } diff --git a/src/java/org/apache/cassandra/streaming/IncomingStreamReader.java b/src/java/org/apache/cassandra/streaming/IncomingStreamReader.java index 3ca6b0e9ac..3f6f9fade9 100644 --- a/src/java/org/apache/cassandra/streaming/IncomingStreamReader.java +++ b/src/java/org/apache/cassandra/streaming/IncomingStreamReader.java @@ -28,6 +28,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.cassandra.net.FileStreamTask; +import org.apache.cassandra.utils.Pair; public class IncomingStreamReader { @@ -58,39 +59,43 @@ public class IncomingStreamReader FileOutputStream fos = new FileOutputStream(pendingFile.getFilename(), true); FileChannel fc = fos.getChannel(); - long bytesRead = 0; + long offset = 0; try { - while (bytesRead < pendingFile.getExpectedBytes()) { - bytesRead += fc.transferFrom(socketChannel, bytesRead, FileStreamTask.CHUNK_SIZE); - pendingFile.update(bytesRead); + Pair section; + while ((section = pendingFile.currentSection()) != null) + { + long length = Math.min(FileStreamTask.CHUNK_SIZE, section.right - section.left); + long bytesRead = fc.transferFrom(socketChannel, offset, length); + // offset in the remote file + pendingFile.update(section.left + bytesRead); + // offset in the local file + offset += bytesRead; } - logger.debug("Receiving stream: finished reading chunk, awaiting more"); } catch (IOException ex) { + logger.debug("Receiving stream: recovering from IO error"); /* Ask the source node to re-stream this file. */ streamStatus.setAction(FileStatus.Action.STREAM); handleFileStatus(remoteAddress.getAddress()); /* Delete the orphaned file. */ File file = new File(pendingFile.getFilename()); file.delete(); - logger.debug("Receiving stream: recovering from IO error"); + /* Reset our state. */ + pendingFile.update(0); throw ex; } finally { + fc.close(); StreamInManager.activeStreams.remove(remoteAddress.getAddress(), pendingFile); } - if (bytesRead == pendingFile.getExpectedBytes()) - { - if (logger.isDebugEnabled()) - logger.debug("Removing stream context " + pendingFile); - fc.close(); - streamStatus.setAction(FileStatus.Action.DELETE); - handleFileStatus(remoteAddress.getAddress()); - } + if (logger.isDebugEnabled()) + logger.debug("Removing stream context " + pendingFile); + streamStatus.setAction(FileStatus.Action.DELETE); + handleFileStatus(remoteAddress.getAddress()); } private void handleFileStatus(InetAddress remoteHost) throws IOException diff --git a/src/java/org/apache/cassandra/streaming/PendingFile.java b/src/java/org/apache/cassandra/streaming/PendingFile.java index 97966b4966..007e6630af 100644 --- a/src/java/org/apache/cassandra/streaming/PendingFile.java +++ b/src/java/org/apache/cassandra/streaming/PendingFile.java @@ -24,17 +24,23 @@ package org.apache.cassandra.streaming; import java.io.DataInputStream; import java.io.DataOutputStream; import java.io.IOException; +import java.util.ArrayList; +import java.util.List; import org.apache.cassandra.io.ICompactSerializer; import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.utils.Pair; -class PendingFile +/** + * Represents portions of a file to be streamed between nodes. + */ +public class PendingFile { private static ICompactSerializer serializer_; static { - serializer_ = new InitiatedFileSerializer(); + serializer_ = new PendingFileSerializer(); } public static ICompactSerializer serializer() @@ -42,16 +48,21 @@ class PendingFile return serializer_; } - private Descriptor desc; - private String component; - private long expectedBytes; + private final Descriptor desc; + private final String component; + private final List> sections; private long ptr; - public PendingFile(Descriptor desc, String component, long expectedBytes) + public PendingFile(Descriptor desc, PendingFile pf) + { + this(desc, pf.component, pf.sections); + } + + public PendingFile(Descriptor desc, String component, List> sections) { this.desc = desc; this.component = component; - this.expectedBytes = expectedBytes; + this.sections = sections; ptr = 0; } @@ -60,9 +71,16 @@ class PendingFile this.ptr = ptr; } - public long getPtr() + /** + * @return The current section of the file, as an (offset,end) pair, or null if nothing left to stream. + */ + public Pair currentSection() { - return ptr; + // linear search for the first appropriate section + for (Pair section : sections) + if (ptr < section.right) + return new Pair(Long.valueOf(Math.max(ptr, section.left)), section.right); + return null; } public String getComponent() @@ -80,11 +98,6 @@ class PendingFile return desc.filenameFor(component); } - public long getExpectedBytes() - { - return expectedBytes; - } - public boolean equals(Object o) { if ( !(o instanceof PendingFile) ) @@ -96,29 +109,36 @@ class PendingFile public int hashCode() { - return toString().hashCode(); + return getFilename().hashCode(); } public String toString() { - return getFilename() + ":" + expectedBytes; + return getFilename() + ":" + ptr + "/" + sections; } - private static class InitiatedFileSerializer implements ICompactSerializer + private static class PendingFileSerializer implements ICompactSerializer { public void serialize(PendingFile sc, DataOutputStream dos) throws IOException { dos.writeUTF(sc.desc.filenameFor(sc.component)); dos.writeUTF(sc.component); - dos.writeLong(sc.expectedBytes); + dos.writeInt(sc.sections.size()); + for (Pair section : sc.sections) + { + dos.writeLong(section.left); dos.writeLong(section.right); + } } public PendingFile deserialize(DataInputStream dis) throws IOException { Descriptor desc = Descriptor.fromFilename(dis.readUTF()); String component = dis.readUTF(); - long expectedBytes = dis.readLong(); - return new PendingFile(desc, component, expectedBytes); + int count = dis.readInt(); + List> sections = new ArrayList>(count); + for (int i = 0; i < count; i++) + sections.add(new Pair(Long.valueOf(dis.readLong()), Long.valueOf(dis.readLong()))); + return new PendingFile(desc, component, sections); } } } diff --git a/src/java/org/apache/cassandra/streaming/StreamFinishedVerbHandler.java b/src/java/org/apache/cassandra/streaming/StreamFinishedVerbHandler.java index 3024e724f6..2c3c39b3d8 100644 --- a/src/java/org/apache/cassandra/streaming/StreamFinishedVerbHandler.java +++ b/src/java/org/apache/cassandra/streaming/StreamFinishedVerbHandler.java @@ -49,7 +49,7 @@ public class StreamFinishedVerbHandler implements IVerbHandler switch (streamStatus.getAction()) { case DELETE: - StreamOutManager.get(message.getFrom()).finishAndStartNext(streamStatus.getFile()); + StreamOutManager.get(message.getFrom()).finishAndStartNext(); break; case STREAM: diff --git a/src/java/org/apache/cassandra/streaming/StreamInitiateMessage.java b/src/java/org/apache/cassandra/streaming/StreamInitiateMessage.java index 459a890bca..9f7577e209 100644 --- a/src/java/org/apache/cassandra/streaming/StreamInitiateMessage.java +++ b/src/java/org/apache/cassandra/streaming/StreamInitiateMessage.java @@ -77,15 +77,9 @@ class StreamInitiateMessage public StreamInitiateMessage deserialize(DataInputStream dis) throws IOException { int size = dis.readInt(); - PendingFile[] pendingFiles = new PendingFile[0]; - if ( size > 0 ) - { - pendingFiles = new PendingFile[size]; - for ( int i = 0; i < size; ++i ) - { - pendingFiles[i] = PendingFile.serializer().deserialize(dis); - } - } + PendingFile[] pendingFiles = new PendingFile[size]; + for (int i = 0; i < size; i++) + pendingFiles[i] = PendingFile.serializer().deserialize(dis); return new StreamInitiateMessage(pendingFiles); } } diff --git a/src/java/org/apache/cassandra/streaming/StreamInitiateVerbHandler.java b/src/java/org/apache/cassandra/streaming/StreamInitiateVerbHandler.java index 22cf48dfc2..6fed8f8e2b 100644 --- a/src/java/org/apache/cassandra/streaming/StreamInitiateVerbHandler.java +++ b/src/java/org/apache/cassandra/streaming/StreamInitiateVerbHandler.java @@ -76,7 +76,7 @@ public class StreamInitiateVerbHandler implements IVerbHandler PendingFile remoteFile = pendingFile.getKey(); PendingFile localFile = pendingFile.getValue(); - FileStatus streamStatus = new FileStatus(remoteFile.getFilename(), remoteFile.getExpectedBytes()); + FileStatus streamStatus = new FileStatus(remoteFile.getFilename()); if (logger.isDebugEnabled()) logger.debug("Preparing to receive stream from " + message.getFrom() + ": " + remoteFile + " -> " + localFile); @@ -114,7 +114,7 @@ public class StreamInitiateVerbHandler implements IVerbHandler Descriptor localdesc = Descriptor.fromFilename(cfStore.getFlushPath()); // add a local file for this component - mapping.put(remote, new PendingFile(localdesc, remote.getComponent(), remote.getExpectedBytes())); + mapping.put(remote, new PendingFile(localdesc, remote)); } return mapping; diff --git a/src/java/org/apache/cassandra/streaming/StreamOut.java b/src/java/org/apache/cassandra/streaming/StreamOut.java index 8d30813197..f9cbe1207f 100644 --- a/src/java/org/apache/cassandra/streaming/StreamOut.java +++ b/src/java/org/apache/cassandra/streaming/StreamOut.java @@ -38,7 +38,7 @@ import org.apache.cassandra.io.sstable.SSTable; import org.apache.cassandra.io.sstable.SSTableReader; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; - +import org.apache.cassandra.utils.Pair; /** * This class handles streaming data from one node to another. @@ -73,7 +73,7 @@ public class StreamOut /* * (1) dump all the memtables to disk. - * (2) anticompaction -- split out the keys in the range specified + * (2) determine the minimal file sections we need to send for the given ranges * (3) transfer the data. */ try @@ -95,9 +95,8 @@ public class StreamOut throw new RuntimeException(e); } } - logger.info("Performing anticompaction ..."); - /* Get the list of files that need to be streamed */ - transferSSTables(target, table.forceAntiCompaction(ranges, target), tableName); // SSTR GC deletes the file when done + // send the matching portion of every sstable in the keyspace + transferSSTables(target, tableName, table.getAllSSTables(), ranges); } catch (IOException e) { @@ -112,20 +111,23 @@ public class StreamOut } /** - * Transfers a group of sstables from a single table to the target endpoint - * and then marks them as ready for local deletion. + * Transfers matching portions of a group of sstables from a single table to the target endpoint. */ - public static void transferSSTables(InetAddress target, List sstables, String table) throws IOException + public static void transferSSTables(InetAddress target, String table, Collection sstables, Collection ranges) throws IOException { - PendingFile[] pendingFiles = new PendingFile[sstables.size()]; + List pending = new ArrayList(); int i = 0; for (SSTableReader sstable : sstables) { Descriptor desc = sstable.getDescriptor(); - long filelen = new File(desc.filenameFor(SSTable.COMPONENT_DATA)).length(); - pendingFiles[i++] = new PendingFile(desc, SSTable.COMPONENT_DATA, filelen); + List> sections = sstable.getPositionsForRanges(ranges); + if (sections.isEmpty()) + continue; + pending.add(new PendingFile(desc, SSTable.COMPONENT_DATA, sections)); } - logger.info("Stream context metadata " + StringUtils.join(pendingFiles, ", " + " " + sstables.size() + " sstables.")); + logger.info("Stream context metadata " + pending + " " + sstables.size() + " sstables."); + + PendingFile[] pendingFiles = pending.toArray(new PendingFile[pending.size()]); StreamOutManager.get(target).addFilesToStream(pendingFiles); StreamInitiateMessage biMessage = new StreamInitiateMessage(pendingFiles); Message message = StreamInitiateMessage.makeStreamInitiateMessage(biMessage); diff --git a/src/java/org/apache/cassandra/streaming/StreamOutManager.java b/src/java/org/apache/cassandra/streaming/StreamOutManager.java index 0bef13d966..ef0c6d9b23 100644 --- a/src/java/org/apache/cassandra/streaming/StreamOutManager.java +++ b/src/java/org/apache/cassandra/streaming/StreamOutManager.java @@ -35,6 +35,7 @@ import java.net.InetAddress; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.io.util.FileUtils; +import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.SimpleCondition; import org.slf4j.Logger; @@ -96,7 +97,6 @@ public class StreamOutManager private final Map fileMap = new HashMap(); private final InetAddress to; - private long totalBytes = 0L; private final SimpleCondition condition = new SimpleCondition(); private StreamOutManager(InetAddress to) @@ -114,41 +114,39 @@ public class StreamOutManager logger.debug("Adding file " + pendingFile.getFilename() + " to be streamed."); files.add(pendingFile); fileMap.put(pendingFile.getFilename(), pendingFile); - totalBytes += pendingFile.getExpectedBytes(); } } + /** + * An (offset,end) pair representing the current section of the file to stream. + */ + public Pair currentSection(String path) + { + return fileMap.get(path).currentSection(); + } + public void update(String path, long pos) { - PendingFile pf = fileMap.get(path); - if (pf != null) - pf.update(pos); + fileMap.get(path).update(pos); } public void startNext() { if (files.size() > 0) { - File file = new File(files.get(0).getFilename()); + PendingFile pf = files.get(0); if (logger.isDebugEnabled()) - logger.debug("Streaming " + file.length() + " length file " + file + " ..."); - MessagingService.instance.stream(file.getAbsolutePath(), to); + logger.debug("Streaming " + pf + " ..."); + MessagingService.instance.stream(pf, to); } } - public void finishAndStartNext(String file) throws IOException + public void finishAndStartNext() throws IOException { - File f = new File(file); - if (logger.isDebugEnabled()) - logger.debug("Deleting file " + file + " after streaming " + f.length() + "/" + totalBytes + " bytes."); - FileUtils.delete(file); PendingFile pf = files.remove(0); - if (pf != null) - fileMap.remove(pf.getFilename()); + fileMap.remove(pf.getFilename()); if (files.size() > 0) - { startNext(); - } else { if (logger.isDebugEnabled()) diff --git a/src/java/org/apache/cassandra/streaming/StreamingService.java b/src/java/org/apache/cassandra/streaming/StreamingService.java index d59518084d..2e9c6c0919 100644 --- a/src/java/org/apache/cassandra/streaming/StreamingService.java +++ b/src/java/org/apache/cassandra/streaming/StreamingService.java @@ -58,7 +58,7 @@ public class StreamingService implements StreamingServiceMBean sb.append(String.format(" %s:\n", source.getHostAddress())); for (PendingFile pf : StreamInManager.getIncomingFiles(source)) { - sb.append(String.format(" %s %d/%d\n", pf.getFilename(), pf.getPtr(), pf.getExpectedBytes())); + sb.append(String.format(" %s\n", pf.toString())); } } sb.append("Sending to:\n"); @@ -67,7 +67,7 @@ public class StreamingService implements StreamingServiceMBean sb.append(String.format(" %s:\n", dest.getHostAddress())); for (PendingFile pf : StreamOutManager.getPendingFiles(dest)) { - sb.append(String.format(" %s %d/%d\n", pf.getFilename(), pf.getPtr(), pf.getExpectedBytes())); + sb.append(String.format(" %s\n", pf.toString())); } } return sb.toString(); @@ -92,7 +92,7 @@ public class StreamingService implements StreamingServiceMBean StreamOutManager manager = StreamOutManager.get(dest); for (PendingFile f : manager.getFiles()) - files.add(String.format("%s %d/%d", f.getFilename(), f.getPtr(), f.getExpectedBytes())); + files.add(String.format("%s", f.toString())); return files; } @@ -108,7 +108,7 @@ public class StreamingService implements StreamingServiceMBean List files = new ArrayList(); for (PendingFile pf : StreamInManager.getIncomingFiles(InetAddress.getByName(host))) { - files.add(String.format("%s: %s %d/%d", pf.getDescriptor().ksname, pf.getFilename(), pf.getPtr(), pf.getExpectedBytes())); + files.add(String.format("%s: %s", pf.getDescriptor().ksname, pf.toString())); } return files; } diff --git a/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java b/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java index 6df0a3d23c..37104a5df7 100644 --- a/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java +++ b/test/unit/org/apache/cassandra/db/ColumnFamilyStoreTest.java @@ -64,7 +64,7 @@ public class ColumnFamilyStoreTest extends CleanupHelper ColumnFamilyStore store = Util.writeColumnFamily(rms); Table table = Table.open("Keyspace1"); - List ssTables = table.getAllSSTablesOnDisk(); + List ssTables = table.getAllSSTables(); assertEquals(1, ssTables.size()); ssTables.get(0).forceFilterFailures(); ColumnFamily cf = store.getColumnFamily(QueryFilter.getIdentityFilter(Util.dk("key2"), new QueryPath("Standard1", null, "Column1".getBytes()))); diff --git a/test/unit/org/apache/cassandra/db/TableTest.java b/test/unit/org/apache/cassandra/db/TableTest.java index f6ef86e5ac..c9d0c1b6ad 100644 --- a/test/unit/org/apache/cassandra/db/TableTest.java +++ b/test/unit/org/apache/cassandra/db/TableTest.java @@ -431,7 +431,7 @@ public class TableTest extends CleanupHelper CompactionManager.instance.submitMajor(cfStore).get(); } SSTableReader sstable = cfStore.getSSTables().iterator().next(); - long position = sstable.getPosition(key); + long position = sstable.getPosition(key, SSTableReader.Operator.EQ); BufferedRandomAccessFile file = new BufferedRandomAccessFile(sstable.getFilename(), "r"); file.seek(position); assert Arrays.equals(FBUtilities.readShortByteArray(file), key.key); diff --git a/test/unit/org/apache/cassandra/io/StreamingTest.java b/test/unit/org/apache/cassandra/io/StreamingTest.java index 0cd2bceed6..7f62b81f1e 100644 --- a/test/unit/org/apache/cassandra/io/StreamingTest.java +++ b/test/unit/org/apache/cassandra/io/StreamingTest.java @@ -26,6 +26,11 @@ import java.util.*; import org.apache.cassandra.CleanupHelper; import org.apache.cassandra.Util; import org.apache.cassandra.db.*; +import org.apache.cassandra.db.filter.QueryFilter; +import org.apache.cassandra.db.filter.QueryPath; +import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; import org.apache.cassandra.io.sstable.SSTableUtils; import org.apache.cassandra.io.sstable.SSTableReader; import org.apache.cassandra.service.StorageService; @@ -46,17 +51,28 @@ public class StreamingTest extends CleanupHelper // write a temporary SSTable, but don't register it Set content = new HashSet(); content.add("key"); + content.add("key2"); + content.add("key3"); SSTableReader sstable = SSTableUtils.writeSSTable(content); String tablename = sstable.getTableName(); String cfname = sstable.getColumnFamilyName(); - // transfer - StreamOut.transferSSTables(LOCAL, Arrays.asList(sstable), tablename); + // transfer the first and last key + IPartitioner p = StorageService.getPartitioner(); + List ranges = new ArrayList(); + ranges.add(new Range(p.getMinimumToken(), p.getToken("key".getBytes()))); + ranges.add(new Range(p.getToken("key2".getBytes()), p.getMinimumToken())); + StreamOut.transferSSTables(LOCAL, tablename, Arrays.asList(sstable), ranges); // confirm that the SSTable was transferred and registered ColumnFamilyStore cfstore = Table.open(tablename).getColumnFamilyStore(cfname); List rows = Util.getRangeSlice(cfstore); - assert rows.size() == 1; + assertEquals(2, rows.size()); assert Arrays.equals(rows.get(0).key.key, "key".getBytes()); + assert Arrays.equals(rows.get(1).key.key, "key3".getBytes()); + + // and that the index and filter were properly recovered + assert null != cfstore.getColumnFamily(QueryFilter.getIdentityFilter(Util.dk("key"), new QueryPath("Standard1"))); + assert null != cfstore.getColumnFamily(QueryFilter.getIdentityFilter(Util.dk("key3"), new QueryPath("Standard1"))); } } diff --git a/test/unit/org/apache/cassandra/io/sstable/LegacySSTableTest.java b/test/unit/org/apache/cassandra/io/sstable/LegacySSTableTest.java index b44b331dff..46de761c56 100644 --- a/test/unit/org/apache/cassandra/io/sstable/LegacySSTableTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/LegacySSTableTest.java @@ -102,7 +102,7 @@ public class LegacySSTableTest extends CleanupHelper for (byte[] key : keys) { // confirm that the bloom filter does not reject any keys - file.seek(reader.getPosition(reader.partitioner.decorateKey(key))); + file.seek(reader.getPosition(reader.partitioner.decorateKey(key), SSTableReader.Operator.EQ)); assert Arrays.equals(key, FBUtilities.readShortByteArray(file)); } } diff --git a/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java b/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java index 99b0fffcb5..9691d5608c 100644 --- a/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java @@ -2,6 +2,9 @@ package org.apache.cassandra.io.sstable; import java.io.IOException; import java.util.concurrent.ExecutionException; +import java.util.ArrayList; +import java.util.Map; +import java.util.List; import org.junit.Test; @@ -9,9 +12,14 @@ import org.apache.cassandra.CleanupHelper; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.*; import org.apache.cassandra.db.filter.QueryPath; +import org.apache.cassandra.dht.Range; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.io.util.BufferedRandomAccessFile; import org.apache.cassandra.io.util.FileDataInput; import org.apache.cassandra.io.util.MmappedSegmentedFile; +import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.FBUtilities; +import org.apache.cassandra.utils.Pair; import org.apache.cassandra.Util; @@ -19,6 +27,50 @@ import static org.junit.Assert.assertEquals; public class SSTableReaderTest extends CleanupHelper { + static Token t(int i) + { + return StorageService.getPartitioner().getToken(String.valueOf(i).getBytes()); + } + + @Test + public void testGetPositionsForRanges() throws IOException, ExecutionException, InterruptedException + { + Table table = Table.open("Keyspace1"); + ColumnFamilyStore store = table.getColumnFamilyStore("Standard2"); + + // insert data and compact to a single sstable + CompactionManager.instance.disableAutoCompaction(); + for (int j = 0; j < 10; j++) + { + byte[] key = String.valueOf(j).getBytes(); + RowMutation rm = new RowMutation("Keyspace1", key); + rm.add(new QueryPath("Standard2", null, "0".getBytes()), new byte[0], new TimestampClock(j)); + rm.apply(); + } + store.forceBlockingFlush(); + CompactionManager.instance.submitMajor(store).get(); + + List ranges = new ArrayList(); + // 1 key + ranges.add(new Range(t(0), t(1))); + // 2 keys + ranges.add(new Range(t(2), t(4))); + // wrapping range from key to end + ranges.add(new Range(t(6), StorageService.getPartitioner().getMinimumToken())); + // empty range (should be ignored) + ranges.add(new Range(t(9), t(91))); + + // confirm that positions increase continuously + SSTableReader sstable = store.getSSTables().iterator().next(); + long previous = -1; + for (Pair section : sstable.getPositionsForRanges(ranges)) + { + assert previous <= section.left : previous + " ! < " + section.left; + assert section.left < section.right : section.left + " ! < " + section.right; + previous = section.right; + } + } + @Test public void testSpannedIndexPositions() throws IOException, ExecutionException, InterruptedException { @@ -53,7 +105,7 @@ public class SSTableReaderTest extends CleanupHelper for (int j = 1; j < 110; j += 2) { DecoratedKey dk = Util.dk(String.valueOf(j)); - assert sstable.getPosition(dk) == -1; + assert sstable.getPosition(dk, SSTableReader.Operator.EQ) == -1; } } } diff --git a/test/unit/org/apache/cassandra/io/sstable/SSTableTest.java b/test/unit/org/apache/cassandra/io/sstable/SSTableTest.java index 20a0812152..926c08743f 100644 --- a/test/unit/org/apache/cassandra/io/sstable/SSTableTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/SSTableTest.java @@ -50,7 +50,7 @@ public class SSTableTest extends CleanupHelper private void verifySingle(SSTableReader sstable, byte[] bytes, byte[] key) throws IOException { BufferedRandomAccessFile file = new BufferedRandomAccessFile(sstable.getFilename(), "r"); - file.seek(sstable.getPosition(sstable.partitioner.decorateKey(key))); + file.seek(sstable.getPosition(sstable.partitioner.decorateKey(key), SSTableReader.Operator.EQ)); assert Arrays.equals(key, FBUtilities.readShortByteArray(file)); int size = (int)SSTableReader.readRowSize(file, sstable.getDescriptor()); byte[] bytes2 = new byte[size]; @@ -82,7 +82,7 @@ public class SSTableTest extends CleanupHelper BufferedRandomAccessFile file = new BufferedRandomAccessFile(sstable.getFilename(), "r"); for (byte[] key : keys) { - file.seek(sstable.getPosition(sstable.partitioner.decorateKey(key))); + file.seek(sstable.getPosition(sstable.partitioner.decorateKey(key), SSTableReader.Operator.EQ)); assert Arrays.equals(key, FBUtilities.readShortByteArray(file)); int size = (int)SSTableReader.readRowSize(file, sstable.getDescriptor()); byte[] bytes2 = new byte[size]; diff --git a/test/unit/org/apache/cassandra/streaming/BootstrapTest.java b/test/unit/org/apache/cassandra/streaming/BootstrapTest.java index 295ed6eec9..eff841b8c2 100644 --- a/test/unit/org/apache/cassandra/streaming/BootstrapTest.java +++ b/test/unit/org/apache/cassandra/streaming/BootstrapTest.java @@ -26,6 +26,9 @@ import java.util.Map; import org.apache.cassandra.SchemaLoader; import org.apache.cassandra.io.sstable.Descriptor; +import org.apache.cassandra.utils.Pair; + +import java.util.Arrays; import org.junit.Test; @@ -35,7 +38,7 @@ public class BootstrapTest extends SchemaLoader public void testGetNewNames() throws IOException { Descriptor desc = Descriptor.fromFilename(new File("Keyspace1", "Standard1-500-Data.db").toString()); - PendingFile[] pendingFiles = new PendingFile[]{ new PendingFile(desc, "Data.db", 100) }; + PendingFile[] pendingFiles = new PendingFile[]{ new PendingFile(desc, "Data.db", Arrays.asList(new Pair(0L, 1L))) }; StreamInitiateVerbHandler bivh = new StreamInitiateVerbHandler(); // map the input (remote) contexts to output (local) contexts