From cd60613ae03642a841c0fd460205bcf78b10249c Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Wed, 1 Jun 2011 15:34:57 +0000 Subject: [PATCH 1/2] fix truncate/compaction race patch by jbellis; reviewed by slebresne for CASSANDRA-2673 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1130191 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../cassandra/db/ColumnFamilyStore.java | 45 ++++++++----------- .../cassandra/db/CompactionManager.java | 24 ++++++++++ src/java/org/apache/cassandra/db/Table.java | 17 ------- .../cassandra/db/TruncateVerbHandler.java | 3 +- 5 files changed, 45 insertions(+), 45 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 2470c7dc2d..9d3f21f51a 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -12,6 +12,7 @@ * close scrub file handles (CASSANDRA-2669) * throttle migration replay (CASSANDRA-2714) * optimize column serializer creation (CASSANDRA-2716) + * fix truncate/compaction race (CASSANDRA-2673) 0.7.6 diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 69ca242599..3458613fd8 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -1887,8 +1887,12 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean */ public Future truncate() throws IOException { - // snapshot will also flush, but we want to truncate the most possible, and anything in a flush written - // after truncateAt won't be truncated. + // We have two goals here: + // - truncate should delete everything written before truncate was invoked + // - but not delete anything that isn't part of the snapshot we create. + // We accomplish this by first flushing manually, then snapshotting, and + // recording the timestamp IN BETWEEN those actions. Any sstables created + // with this timestamp or greater time, will not be marked for delete. try { forceBlockingFlush(); @@ -1897,33 +1901,20 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean { throw new RuntimeException(e); } - - final long truncatedAt = System.currentTimeMillis(); + // sleep a little to make sure that our truncatedAt comes after any sstable + // that was part of the flushed we forced; otherwise on a tie, it won't get deleted. + try + { + Thread.sleep(100); + } + catch (InterruptedException e) + { + throw new AssertionError(e); + } + long truncatedAt = System.currentTimeMillis(); snapshot(Table.getTimestampedSnapshotName("before-truncate")); - Runnable runnable = new WrappedRunnable() - { - public void runMayThrow() throws InterruptedException, IOException - { - // putting markCompacted on the commitlogUpdater thread ensures it will run - // after any compactions that were in progress when truncate was called, are finished - for (ColumnFamilyStore cfs : concatWithIndexes()) - { - List truncatedSSTables = new ArrayList(); - for (SSTableReader sstable : cfs.getSSTables()) - { - if (!sstable.newSince(truncatedAt)) - truncatedSSTables.add(sstable); - } - cfs.markCompacted(truncatedSSTables); - } - - // Invalidate row cache - invalidateRowCache(); - } - }; - - return postFlushExecutor.submit(runnable); + return CompactionManager.instance.submitTruncate(this, truncatedAt); } // if this errors out, we are in a world of hurt. diff --git a/src/java/org/apache/cassandra/db/CompactionManager.java b/src/java/org/apache/cassandra/db/CompactionManager.java index 56fa1c365f..c3420a051b 100644 --- a/src/java/org/apache/cassandra/db/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/CompactionManager.java @@ -977,6 +977,30 @@ public class CompactionManager implements CompactionManagerMBean return executor.submit(runnable); } + public Future submitTruncate(final ColumnFamilyStore main, final long truncatedAt) + { + Runnable runnable = new WrappedRunnable() + { + public void runMayThrow() throws InterruptedException, IOException + { + for (ColumnFamilyStore cfs : main.concatWithIndexes()) + { + List truncatedSSTables = new ArrayList(); + for (SSTableReader sstable : cfs.getSSTables()) + { + if (!sstable.newSince(truncatedAt)) + truncatedSSTables.add(sstable); + } + cfs.markCompacted(truncatedSSTables); + } + + main.invalidateRowCache(); + } + }; + + return executor.submit(runnable); + } + private static int getDefaultGcBefore(ColumnFamilyStore cfs) { return (int) (System.currentTimeMillis() / 1000) - cfs.metadata.getGcGraceSeconds(); diff --git a/src/java/org/apache/cassandra/db/Table.java b/src/java/org/apache/cassandra/db/Table.java index f6bb0096b6..a9a7546b5c 100644 --- a/src/java/org/apache/cassandra/db/Table.java +++ b/src/java/org/apache/cassandra/db/Table.java @@ -677,23 +677,6 @@ public class Table return Iterables.transform(DatabaseDescriptor.getTables(), transformer); } - /** - * Performs a synchronous truncate operation, effectively deleting all data - * from the column family cfname - * @param cfname - * @throws IOException - * @throws ExecutionException - * @throws InterruptedException - */ - public void truncate(String cfname) throws InterruptedException, ExecutionException, IOException - { - logger.debug("Truncating..."); - ColumnFamilyStore cfs = getColumnFamilyStore(cfname); - // truncate, blocking - cfs.truncate().get(); - logger.debug("Truncation done."); - } - @Override public String toString() { return getClass().getSimpleName() + "(name='" + name + "')"; diff --git a/src/java/org/apache/cassandra/db/TruncateVerbHandler.java b/src/java/org/apache/cassandra/db/TruncateVerbHandler.java index cebf0bd5d2..bc6cd46fb8 100644 --- a/src/java/org/apache/cassandra/db/TruncateVerbHandler.java +++ b/src/java/org/apache/cassandra/db/TruncateVerbHandler.java @@ -52,7 +52,8 @@ public class TruncateVerbHandler implements IVerbHandler try { - Table.open(t.keyspace).truncate(t.columnFamily); + ColumnFamilyStore cfs = Table.open(t.keyspace).getColumnFamilyStore(t.columnFamily); + cfs.truncate().get(); } catch (IOException e) { From a7bc6ee07cbe6f2952a60d47786d6b152cddfbde Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Jun 2011 23:44:31 +0000 Subject: [PATCH 2/2] merge from 0.6 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1131292 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 5 +++++ .../apache/cassandra/net/IncomingTcpConnection.java | 11 +++++++++-- 2 files changed, 14 insertions(+), 2 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 9d3f21f51a..2664a1aaa3 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -13,6 +13,8 @@ * throttle migration replay (CASSANDRA-2714) * optimize column serializer creation (CASSANDRA-2716) * fix truncate/compaction race (CASSANDRA-2673) + * workaround large resultsets causing large allocation retention + by nio sockets (CASSANDRA-2654) 0.7.6 @@ -54,6 +56,9 @@ * reduce contention on Table.flusherLock (CASSANDRA-1954) * try harder to detect failures during streaming, cleaning up temporary files more reliably (CASSANDRA-2088) + + +0.6.13 * shut down server for OOM on a Thrift thread (CASSANDRA-2269) * fix tombstone handling in repair and sstable2json (CASSANDRA-2279) * preserve version when streaming data from old sstables (CASSANDRA-2283) diff --git a/src/java/org/apache/cassandra/net/IncomingTcpConnection.java b/src/java/org/apache/cassandra/net/IncomingTcpConnection.java index 52c61e9619..cfe66dd700 100644 --- a/src/java/org/apache/cassandra/net/IncomingTcpConnection.java +++ b/src/java/org/apache/cassandra/net/IncomingTcpConnection.java @@ -35,6 +35,8 @@ public class IncomingTcpConnection extends Thread { private static Logger logger = LoggerFactory.getLogger(IncomingTcpConnection.class); + private static final int CHUNK_SIZE = 1024 * 1024; + private Socket socket; public IncomingTcpConnection(Socket socket) @@ -95,8 +97,13 @@ public class IncomingTcpConnection extends Thread { int size = input.readInt(); byte[] contentBytes = new byte[size]; - input.readFully(contentBytes); - + // readFully allocates a direct buffer the size of the chunk it is asked to read, + // so we cap that at CHUNK_SIZE. See https://issues.apache.org/jira/browse/CASSANDRA-2654 + int remainder = size % CHUNK_SIZE; + for (int offset = 0; offset < size - remainder; offset += CHUNK_SIZE) + input.readFully(contentBytes, offset, CHUNK_SIZE); + input.readFully(contentBytes, size - remainder, remainder); + if (version > MessagingService.version_) logger.info("Received connection from newer protocol version. Ignorning message."); else