From 10786e72571ebe9382970834c95301de661bdf5d Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Tue, 27 Sep 2011 20:36:59 +0000 Subject: [PATCH 01/15] merge from 0.8 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1176605 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 + .../apache/cassandra/utils/FBUtilities.java | 10 +++-- .../cassandra/db/marshal/BytesTypeTest.java | 40 +++++++++++++++++++ 3 files changed, 48 insertions(+), 4 deletions(-) create mode 100644 test/unit/org/apache/cassandra/db/marshal/BytesTypeTest.java diff --git a/CHANGES.txt b/CHANGES.txt index 54c0e501e6..c3fdec7c9c 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -4,6 +4,8 @@ * test for NUMA policy support as well as numactl presence (CASSANDRA-3245) * Fix FD leak when internode encryption is enabled (CASSANDRA-3257) * Remove incorrect assertion in mergeIterator (CASSANDRA-3260) + * FBUtilities.hexToBytes(String) to throw NumberFormatException when string + contains non-hex characters (CASSANDRA-3231) 1.0.0-rc1 diff --git a/src/java/org/apache/cassandra/utils/FBUtilities.java b/src/java/org/apache/cassandra/utils/FBUtilities.java index 4a8aa71a6f..3500a33eba 100644 --- a/src/java/org/apache/cassandra/utils/FBUtilities.java +++ b/src/java/org/apache/cassandra/utils/FBUtilities.java @@ -47,8 +47,6 @@ import org.apache.cassandra.concurrent.CreationTimeAwareFuture; import org.apache.cassandra.config.ConfigurationException; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.DecoratedKey; -import org.apache.cassandra.db.marshal.AbstractType; -import org.apache.cassandra.db.marshal.TypeParser; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; @@ -406,10 +404,14 @@ public class FBUtilities { if (str.length() % 2 == 1) str = "0" + str; - byte[] bytes = new byte[str.length()/2]; + byte[] bytes = new byte[str.length() / 2]; for (int i = 0; i < bytes.length; i++) { - bytes[i] = (byte)((charToByte[str.charAt(i * 2)] << 4) | charToByte[str.charAt(i*2 + 1)]); + byte halfByte1 = charToByte[str.charAt(i * 2)]; + byte halfByte2 = charToByte[str.charAt(i * 2 + 1)]; + if (halfByte1 == -1 || halfByte2 == -1) + throw new NumberFormatException("Non-hex characters in " + str); + bytes[i] = (byte)((halfByte1 << 4) | halfByte2); } return bytes; } diff --git a/test/unit/org/apache/cassandra/db/marshal/BytesTypeTest.java b/test/unit/org/apache/cassandra/db/marshal/BytesTypeTest.java new file mode 100644 index 0000000000..4d96b76b1e --- /dev/null +++ b/test/unit/org/apache/cassandra/db/marshal/BytesTypeTest.java @@ -0,0 +1,40 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ +package org.apache.cassandra.db.marshal; + +import org.junit.Test; + +public class BytesTypeTest +{ + private static final String INVALID_HEX = "33AG45F"; // Invalid (has a G) + private static final String VALID_HEX = "33A45F"; + + @Test (expected = MarshalException.class) + public void testFromStringWithInvalidString() + { + BytesType.instance.fromString(INVALID_HEX); + } + + @Test + public void testFromStringWithValidString() + { + BytesType.instance.fromString(VALID_HEX); + } +} From eb3ab6e554ce498a4eeabf940f48b9e163847d0c Mon Sep 17 00:00:00 2001 From: Eric Evans Date: Tue, 27 Sep 2011 22:06:50 +0000 Subject: [PATCH 02/15] update for missed term rename s/bytea/blob/ Patch by eevans; reviewed by jbellis for CASSANDRA-3266 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1176639 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/cql/Cql.g | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/cql/Cql.g b/src/java/org/apache/cassandra/cql/Cql.g index 22fed4fedd..9b48e32955 100644 --- a/src/java/org/apache/cassandra/cql/Cql.g +++ b/src/java/org/apache/cassandra/cql/Cql.g @@ -436,7 +436,7 @@ dropColumnFamilyStatement returns [String cfam] ; comparatorType - : 'bytea' | 'ascii' | 'text' | 'varchar' | 'int' | 'varint' | 'bigint' | 'uuid' | 'counter' | 'boolean' | 'date' | 'float' | 'double' | 'decimal' + : 'blob' | 'ascii' | 'text' | 'varchar' | 'int' | 'varint' | 'bigint' | 'uuid' | 'counter' | 'boolean' | 'date' | 'float' | 'double' | 'decimal' ; term returns [Term item] From f26dffa72c4f7e7a5b15cc2978b12072ff780754 Mon Sep 17 00:00:00 2001 From: Eric Evans Date: Tue, 27 Sep 2011 22:15:04 +0000 Subject: [PATCH 03/15] accept numericals for strategy_options key Patch by Pavel Yaskevich; reviewed by eevans for CASSANDRA-3239 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1176644 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/cql/Cql.g | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/cql/Cql.g b/src/java/org/apache/cassandra/cql/Cql.g index 9b48e32955..c52b27d56a 100644 --- a/src/java/org/apache/cassandra/cql/Cql.g +++ b/src/java/org/apache/cassandra/cql/Cql.g @@ -600,7 +600,7 @@ IDENT ; COMPIDENT - : IDENT ( ':' IDENT)* + : IDENT ( ':' (IDENT | INTEGER))* ; UUID From a6e5b6dd33f5ceb6c73b33bf69d52c6ca7d0bcf9 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Wed, 28 Sep 2011 17:03:01 +0000 Subject: [PATCH 04/15] merge from 0.8 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1176961 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 5 +++ .../locator/AbstractEndpointSnitch.java | 34 +++++++++++++++---- .../AbstractNetworkTopologySnitch.java | 29 ---------------- .../locator/DynamicEndpointSnitch.java | 9 ++--- .../cassandra/locator/SimpleSnitch.java | 14 +++++--- 5 files changed, 44 insertions(+), 47 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index c3fdec7c9c..ba0c0d908e 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -6,6 +6,8 @@ * Remove incorrect assertion in mergeIterator (CASSANDRA-3260) * FBUtilities.hexToBytes(String) to throw NumberFormatException when string contains non-hex characters (CASSANDRA-3231) + * Keep SimpleSnitch proximity ordering unchanged from what the Strategy + generates, as intended (CASSANDRA-3262) 1.0.0-rc1 @@ -112,6 +114,9 @@ * make memtable throughput and column count thresholds no-ops (CASSANDRA-2449) +======= + +>>>>>>> .merge-right.r1176712 0.8.6 * revert CASSANDRA-2388 * change TokenRange.endpoints back to listen/broadcast address to match diff --git a/src/java/org/apache/cassandra/locator/AbstractEndpointSnitch.java b/src/java/org/apache/cassandra/locator/AbstractEndpointSnitch.java index a418fc2bd0..c027993309 100644 --- a/src/java/org/apache/cassandra/locator/AbstractEndpointSnitch.java +++ b/src/java/org/apache/cassandra/locator/AbstractEndpointSnitch.java @@ -20,17 +20,39 @@ package org.apache.cassandra.locator; import java.net.InetAddress; -import java.util.Collection; -import java.util.List; +import java.util.*; public abstract class AbstractEndpointSnitch implements IEndpointSnitch { - public abstract List getSortedListByProximity(InetAddress address, Collection unsortedAddress); - public abstract void sortByProximity(InetAddress address, List addresses); + public abstract int compareEndpoints(InetAddress target, InetAddress a1, InetAddress a2); - public int compareEndpoints(InetAddress target, InetAddress a1, InetAddress a2) + /** + * Sorts the Collection of node addresses by proximity to the given address + * @param address the address to sort by proximity to + * @param unsortedAddress the nodes to sort + * @return a new sorted List + */ + public List getSortedListByProximity(InetAddress address, Collection unsortedAddress) { - return a1.getHostAddress().compareTo(a2.getHostAddress()); + List preferred = new ArrayList(unsortedAddress); + sortByProximity(address, preferred); + return preferred; + } + + /** + * Sorts the List of node addresses, in-place, by proximity to the given address + * @param address the address to sort the proximity by + * @param addresses the nodes to sort + */ + public void sortByProximity(final InetAddress address, List addresses) + { + Collections.sort(addresses, new Comparator() + { + public int compare(InetAddress a1, InetAddress a2) + { + return compareEndpoints(address, a1, a2); + } + }); } public void gossiperStarting() diff --git a/src/java/org/apache/cassandra/locator/AbstractNetworkTopologySnitch.java b/src/java/org/apache/cassandra/locator/AbstractNetworkTopologySnitch.java index f6530a69c1..be81b82d11 100644 --- a/src/java/org/apache/cassandra/locator/AbstractNetworkTopologySnitch.java +++ b/src/java/org/apache/cassandra/locator/AbstractNetworkTopologySnitch.java @@ -47,35 +47,6 @@ public abstract class AbstractNetworkTopologySnitch extends AbstractEndpointSnit */ abstract public String getDatacenter(InetAddress endpoint); - /** - * Sorts the Collection of node addresses by proximity to the given address - * @param address the address to sort by proximity to - * @param addresses the nodes to sort - * @return a new sorted List - */ - public List getSortedListByProximity(final InetAddress address, Collection addresses) - { - List preferred = new ArrayList(addresses); - sortByProximity(address, preferred); - return preferred; - } - - /** - * Sorts the List of node addresses by proximity to the given address - * @param address the address to sort the proximity by - * @param addresses the nodes to sort - */ - public void sortByProximity(final InetAddress address, List addresses) - { - Collections.sort(addresses, new Comparator() - { - public int compare(InetAddress a1, InetAddress a2) - { - return compareEndpoints(address, a1, a2); - } - }); - } - public int compareEndpoints(InetAddress address, InetAddress a1, InetAddress a2) { if (address.equals(a1) && !address.equals(a2)) diff --git a/src/java/org/apache/cassandra/locator/DynamicEndpointSnitch.java b/src/java/org/apache/cassandra/locator/DynamicEndpointSnitch.java index a845b30e65..264e554843 100644 --- a/src/java/org/apache/cassandra/locator/DynamicEndpointSnitch.java +++ b/src/java/org/apache/cassandra/locator/DynamicEndpointSnitch.java @@ -129,6 +129,7 @@ public class DynamicEndpointSnitch extends AbstractEndpointSnitch implements ILa return list; } + @Override public void sortByProximity(final InetAddress address, List addresses) { assert address.equals(FBUtilities.getBroadcastAddress()); // we only know about ourself @@ -144,13 +145,7 @@ public class DynamicEndpointSnitch extends AbstractEndpointSnitch implements ILa private void sortByProximityWithScore(final InetAddress address, List addresses) { - Collections.sort(addresses, new Comparator() - { - public int compare(InetAddress a1, InetAddress a2) - { - return compareEndpoints(address, a1, a2); - } - }); + super.sortByProximity(address, addresses); } private void sortByProximityWithBadness(final InetAddress address, List addresses) diff --git a/src/java/org/apache/cassandra/locator/SimpleSnitch.java b/src/java/org/apache/cassandra/locator/SimpleSnitch.java index 929aa9828a..c47888c200 100644 --- a/src/java/org/apache/cassandra/locator/SimpleSnitch.java +++ b/src/java/org/apache/cassandra/locator/SimpleSnitch.java @@ -39,13 +39,17 @@ public class SimpleSnitch extends AbstractEndpointSnitch { return "datacenter1"; } - - public List getSortedListByProximity(final InetAddress address, Collection addresses) - { - return new ArrayList(addresses); - } + @Override public void sortByProximity(final InetAddress address, List addresses) { + // Optimization to avoid walking the list + } + + public int compareEndpoints(InetAddress target, InetAddress a1, InetAddress a2) + { + // Making all endpoints equal ensures we won't change the original ordering (since + // Collections.sort is guaranteed to be stable) + return 0; } } From c56760338ebb9f5a283b0e32a9d98ea81e273627 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 29 Sep 2011 01:02:09 +0000 Subject: [PATCH 05/15] fix counter entry in jdbc TypesMap patch by Corey Hulen; reviewed by jbellis for CASSANDRA-3268 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177137 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + src/java/org/apache/cassandra/cql/jdbc/TypesMap.java | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/CHANGES.txt b/CHANGES.txt index ba0c0d908e..115af0c8bb 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -8,6 +8,7 @@ contains non-hex characters (CASSANDRA-3231) * Keep SimpleSnitch proximity ordering unchanged from what the Strategy generates, as intended (CASSANDRA-3262) + * fix counter entry in jdbc TypesMap (CASSANDRA-3268) 1.0.0-rc1 diff --git a/src/java/org/apache/cassandra/cql/jdbc/TypesMap.java b/src/java/org/apache/cassandra/cql/jdbc/TypesMap.java index e1d649cdf8..3d99c687b9 100644 --- a/src/java/org/apache/cassandra/cql/jdbc/TypesMap.java +++ b/src/java/org/apache/cassandra/cql/jdbc/TypesMap.java @@ -33,7 +33,7 @@ public class TypesMap map.put("org.apache.cassandra.db.marshal.AsciiType", JdbcAscii.instance); map.put("org.apache.cassandra.db.marshal.BooleanType", JdbcBoolean.instance); map.put("org.apache.cassandra.db.marshal.BytesType", JdbcBytes.instance); - map.put("org.apache.cassandra.db.marshal.ColumnCounterType", JdbcCounterColumn.instance); + map.put("org.apache.cassandra.db.marshal.CounterColumnType", JdbcCounterColumn.instance); map.put("org.apache.cassandra.db.marshal.DateType", JdbcDate.instance); map.put("org.apache.cassandra.db.marshal.DecimalType", JdbcDecimal.instance); map.put("org.apache.cassandra.db.marshal.DoubleType", JdbcDouble.instance); From 7954d69001dfb4205c45dde88d2f7b3e9457f31d Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 29 Sep 2011 17:03:14 +0000 Subject: [PATCH 06/15] fix full queue scenario for ParallelCompactionIterator patch by jbellis; reviewed by slebresne for CASSANDRA-3270 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177365 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../compaction/ParallelCompactionIterable.java | 18 +++++------------- 2 files changed, 6 insertions(+), 13 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 115af0c8bb..e1abc76f82 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -9,6 +9,7 @@ * Keep SimpleSnitch proximity ordering unchanged from what the Strategy generates, as intended (CASSANDRA-3262) * fix counter entry in jdbc TypesMap (CASSANDRA-3268) + * fix full queue scenario for ParallelCompactionIterator (CASSANDRA-3270) 1.0.0-rc1 diff --git a/src/java/org/apache/cassandra/db/compaction/ParallelCompactionIterable.java b/src/java/org/apache/cassandra/db/compaction/ParallelCompactionIterable.java index e9bf574d40..a557e408d1 100644 --- a/src/java/org/apache/cassandra/db/compaction/ParallelCompactionIterable.java +++ b/src/java/org/apache/cassandra/db/compaction/ParallelCompactionIterable.java @@ -144,21 +144,13 @@ public class ParallelCompactionIterable extends AbstractCompactionIterable private class Reducer extends MergeIterator.Reducer { private final List rows = new ArrayList(); - - private final ThreadPoolExecutor executor; private int row = 0; - private Reducer() - { - super(); - executor = new ThreadPoolExecutor(Runtime.getRuntime().availableProcessors(), - Runtime.getRuntime().availableProcessors(), - Integer.MAX_VALUE, - TimeUnit.MILLISECONDS, - new SynchronousQueue(), - new NamedThreadFactory("CompactionReducer")); - executor.setRejectedExecutionHandler(DebuggableThreadPoolExecutor.blockingExecutionHandler); - } + private final ThreadPoolExecutor executor = new DebuggableThreadPoolExecutor(Runtime.getRuntime().availableProcessors(), + Integer.MAX_VALUE, + TimeUnit.MILLISECONDS, + new SynchronousQueue(), + new NamedThreadFactory("CompactionReducer")); public void reduce(RowContainer current) { From 5e4d3055e9fafd052638099d77b102d81ad67317 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Fri, 30 Sep 2011 15:48:48 +0000 Subject: [PATCH 07/15] Fix bootstrap process patch by slebresne; reviewed by jbellis for CASSANDRA-3285 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177706 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../apache/cassandra/service/StorageService.java | 14 +++++++------- 2 files changed, 8 insertions(+), 7 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index e1abc76f82..422b75098d 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -10,6 +10,7 @@ generates, as intended (CASSANDRA-3262) * fix counter entry in jdbc TypesMap (CASSANDRA-3268) * fix full queue scenario for ParallelCompactionIterator (CASSANDRA-3270) + * fix bootstrap process (CASSANDRA-3285) 1.0.0-rc1 diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index a3192411f6..f0cdce3a7f 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -497,7 +497,7 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe if (DatabaseDescriptor.isAutoBootstrap() && !(SystemTable.isBootstrapped() || DatabaseDescriptor.getSeeds().contains(FBUtilities.getBroadcastAddress()) - || Schema.instance.getNonSystemTables().isEmpty())) + || !Schema.instance.getNonSystemTables().isEmpty())) { setMode("Joining: waiting for ring and schema information", true); try @@ -565,13 +565,13 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe { logger_.info("Using saved token " + token); } - - // start participating in the ring. - SystemTable.setBootstrapped(true); - setToken(token); - logger_.info("Bootstrap/Replace/Move completed! Now serving reads."); - assert tokenMetadata_.sortedTokens().size() > 0; } + + // start participating in the ring. + SystemTable.setBootstrapped(true); + setToken(token); + logger_.info("Bootstrap/Replace/Move completed! Now serving reads."); + assert tokenMetadata_.sortedTokens().size() > 0; } public synchronized void joinRing() throws IOException, org.apache.cassandra.config.ConfigurationException From 75ffb07bbe96924ec5df277c2457ab12ec32ea8f Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 30 Sep 2011 16:29:09 +0000 Subject: [PATCH 08/15] update CHANGES git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177725 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 422b75098d..0dffa624f7 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,6 +1,6 @@ 1.0.0-final - * Log a miningfull warning when a node receive a message for a repair session - that don't exist anymore (CASSANDRA-3256) + * Log a meaningful warning when a node receives a message for a repair session + that doesn't exist anymore (CASSANDRA-3256) * test for NUMA policy support as well as numactl presence (CASSANDRA-3245) * Fix FD leak when internode encryption is enabled (CASSANDRA-3257) * Remove incorrect assertion in mergeIterator (CASSANDRA-3260) From 721664f6806ffe625111b5c010e82982d6839b00 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Fri, 30 Sep 2011 16:31:45 +0000 Subject: [PATCH 09/15] Update version for 1.0.0-rc2 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177726 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 +- build.xml | 2 +- debian/changelog | 6 ++++++ 3 files changed, 8 insertions(+), 2 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 0dffa624f7..3e4ef70611 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,4 @@ -1.0.0-final +1.0.0-rc2 * Log a meaningful warning when a node receives a message for a repair session that doesn't exist anymore (CASSANDRA-3256) * test for NUMA policy support as well as numactl presence (CASSANDRA-3245) diff --git a/build.xml b/build.xml index 7f73700ed2..15ebee1be1 100644 --- a/build.xml +++ b/build.xml @@ -25,7 +25,7 @@ - + diff --git a/debian/changelog b/debian/changelog index cb7c7b7989..50fd43b305 100644 --- a/debian/changelog +++ b/debian/changelog @@ -1,3 +1,9 @@ +cassandra (1.0.0~rc2) unstable; urgency=low + + * New release candidate + + -- Sylvain Lebresne Fri, 30 Sep 2011 18:29:44 +0200 + cassandra (1.0.0~rc1) unstable; urgency=low * New release candidate From 1e7cde5e75ce18cbd88c9119fef58d97f83cc74e Mon Sep 17 00:00:00 2001 From: T Jake Luciani Date: Fri, 30 Sep 2011 17:18:50 +0000 Subject: [PATCH 10/15] Thrift sockets are not properly buffered Patch my tjake; reviewed by jfarrell for CASSANDRA-3261 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0@1177737 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../cassandra/thrift/TCustomServerSocket.java | 91 +++++++- .../cassandra/thrift/TCustomSocket.java | 211 ++++++++++++++++++ 3 files changed, 293 insertions(+), 10 deletions(-) create mode 100644 src/java/org/apache/cassandra/thrift/TCustomSocket.java diff --git a/CHANGES.txt b/CHANGES.txt index 3674553466..219bcb60d7 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,6 @@ 1.0.1 * describe_ring should include datacenter/topology information (CASSANDRA-2882) + * Thrift sockets are not properly buffered (CASSANDRA-3261) 1.0.0-final diff --git a/src/java/org/apache/cassandra/thrift/TCustomServerSocket.java b/src/java/org/apache/cassandra/thrift/TCustomServerSocket.java index d4e273ad9c..38577f1c73 100644 --- a/src/java/org/apache/cassandra/thrift/TCustomServerSocket.java +++ b/src/java/org/apache/cassandra/thrift/TCustomServerSocket.java @@ -1,4 +1,5 @@ package org.apache.cassandra.thrift; + /* * * Licensed to the Apache Software Foundation (ASF) under one @@ -20,8 +21,9 @@ package org.apache.cassandra.thrift; * */ - +import java.io.IOException; import java.net.InetSocketAddress; +import java.net.ServerSocket; import java.net.Socket; import java.net.SocketException; @@ -29,44 +31,79 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.apache.thrift.transport.TServerSocket; +import org.apache.thrift.transport.TServerTransport; import org.apache.thrift.transport.TSocket; import org.apache.thrift.transport.TTransportException; /** - * Extends Thrift's TServerSocket to allow customization of various desirable - * TCP properties. + * Extends Thrift's TServerSocket to allow customization of various desirable TCP properties. */ -public class TCustomServerSocket extends TServerSocket +public class TCustomServerSocket extends TServerTransport { private static final Logger logger = LoggerFactory.getLogger(TCustomServerSocket.class); + /** + * Underlying serversocket object + */ + private ServerSocket serverSocket_ = null; + private final boolean keepAlive; private final Integer sendBufferSize; private final Integer recvBufferSize; /** * Allows fine-tuning of the server socket including keep-alive, reuse of addresses, send and receive buffer sizes. + * * @param bindAddr * @param keepAlive * @param sendBufferSize * @param recvBufferSize * @throws TTransportException */ - public TCustomServerSocket(InetSocketAddress bindAddr, boolean keepAlive, Integer sendBufferSize, Integer recvBufferSize) - throws TTransportException + public TCustomServerSocket(InetSocketAddress bindAddr, boolean keepAlive, Integer sendBufferSize, + Integer recvBufferSize) + throws TTransportException { - super(bindAddr); + try + { + // Make server socket + serverSocket_ = new ServerSocket(); + // Prevent 2MSL delay problem on server restarts + serverSocket_.setReuseAddress(true); + // Bind to listening port + serverSocket_.bind(bindAddr); + } + catch (IOException ioe) + { + serverSocket_ = null; + throw new TTransportException("Could not create ServerSocket on address " + bindAddr.toString() + "."); + } + this.keepAlive = keepAlive; this.sendBufferSize = sendBufferSize; this.recvBufferSize = recvBufferSize; } @Override - protected TSocket acceptImpl() throws TTransportException + protected TCustomSocket acceptImpl() throws TTransportException { - TSocket tsocket = super.acceptImpl(); - Socket socket = tsocket.getSocket(); + + if (serverSocket_ == null) + throw new TTransportException(TTransportException.NOT_OPEN, "No underlying server socket."); + + TCustomSocket tsocket = null; + Socket socket = null; + try + { + socket = serverSocket_.accept(); + tsocket = new TCustomSocket(socket); + tsocket.setTimeout(0); + } + catch (IOException iox) + { + throw new TTransportException(iox); + } try { @@ -103,4 +140,38 @@ public class TCustomServerSocket extends TServerSocket return tsocket; } + + @Override + public void listen() throws TTransportException + { + // Make sure not to block on accept + if (serverSocket_ != null) + { + try + { + serverSocket_.setSoTimeout(0); + } + catch (SocketException sx) + { + logger.error("Could not set socket timeout.", sx); + } + } + } + + @Override + public void close() + { + if (serverSocket_ != null) + { + try + { + serverSocket_.close(); + } + catch (IOException iox) + { + logger.warn("Could not close server socket.", iox); + } + serverSocket_ = null; + } + } } diff --git a/src/java/org/apache/cassandra/thrift/TCustomSocket.java b/src/java/org/apache/cassandra/thrift/TCustomSocket.java new file mode 100644 index 0000000000..90e571362d --- /dev/null +++ b/src/java/org/apache/cassandra/thrift/TCustomSocket.java @@ -0,0 +1,211 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.cassandra.thrift; + + +import java.io.BufferedInputStream; +import java.io.BufferedOutputStream; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.net.Socket; +import java.net.SocketException; + +import org.apache.thrift.transport.TIOStreamTransport; +import org.apache.thrift.transport.TTransportException; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Socket implementation of the TTransport interface. + * + * Adds socket buffering + * + */ +public class TCustomSocket extends TIOStreamTransport { + + private static final Logger LOGGER = LoggerFactory.getLogger(TCustomSocket.class.getName()); + + /** + * Wrapped Socket object + */ + private Socket socket_ = null; + + /** + * Remote host + */ + private String host_ = null; + + /** + * Remote port + */ + private int port_ = 0; + + /** + * Socket timeout + */ + private int timeout_ = 0; + + /** + * Constructor that takes an already created socket. + * + * @param socket Already created socket object + * @throws TTransportException if there is an error setting up the streams + */ + public TCustomSocket(Socket socket) throws TTransportException { + socket_ = socket; + try { + socket_.setSoLinger(false, 0); + socket_.setTcpNoDelay(true); + } catch (SocketException sx) { + LOGGER.warn("Could not configure socket.", sx); + } + + if (isOpen()) { + try { + inputStream_ = new BufferedInputStream(socket_.getInputStream(), 1024); + outputStream_ = new BufferedOutputStream(socket_.getOutputStream(), 1024); + } catch (IOException iox) { + close(); + throw new TTransportException(TTransportException.NOT_OPEN, iox); + } + } + } + + /** + * Creates a new unconnected socket that will connect to the given host + * on the given port. + * + * @param host Remote host + * @param port Remote port + */ + public TCustomSocket(String host, int port) { + this(host, port, 0); + } + + /** + * Creates a new unconnected socket that will connect to the given host + * on the given port. + * + * @param host Remote host + * @param port Remote port + * @param timeout Socket timeout + */ + public TCustomSocket(String host, int port, int timeout) { + host_ = host; + port_ = port; + timeout_ = timeout; + initSocket(); + } + + /** + * Initializes the socket object + */ + private void initSocket() { + socket_ = new Socket(); + try { + socket_.setSoLinger(false, 0); + socket_.setTcpNoDelay(true); + socket_.setSoTimeout(timeout_); + } catch (SocketException sx) { + LOGGER.error("Could not configure socket.", sx); + } + } + + /** + * Sets the socket timeout + * + * @param timeout Milliseconds timeout + */ + public void setTimeout(int timeout) { + timeout_ = timeout; + try { + socket_.setSoTimeout(timeout); + } catch (SocketException sx) { + LOGGER.warn("Could not set socket timeout.", sx); + } + } + + /** + * Returns a reference to the underlying socket. + */ + public Socket getSocket() { + if (socket_ == null) { + initSocket(); + } + return socket_; + } + + /** + * Checks whether the socket is connected. + */ + public boolean isOpen() { + if (socket_ == null) { + return false; + } + return socket_.isConnected(); + } + + /** + * Connects the socket, creating a new socket object if necessary. + */ + public void open() throws TTransportException { + if (isOpen()) { + throw new TTransportException(TTransportException.ALREADY_OPEN, "Socket already connected."); + } + + if (host_.length() == 0) { + throw new TTransportException(TTransportException.NOT_OPEN, "Cannot open null host."); + } + if (port_ <= 0) { + throw new TTransportException(TTransportException.NOT_OPEN, "Cannot open without port."); + } + + if (socket_ == null) { + initSocket(); + } + + try { + socket_.connect(new InetSocketAddress(host_, port_), timeout_); + inputStream_ = new BufferedInputStream(socket_.getInputStream(), 1024); + outputStream_ = new BufferedOutputStream(socket_.getOutputStream(), 1024); + } catch (IOException iox) { + close(); + throw new TTransportException(TTransportException.NOT_OPEN, iox); + } + } + + /** + * Closes the socket. + */ + public void close() { + // Close the underlying streams + super.close(); + + // Close the socket + if (socket_ != null) { + try { + socket_.close(); + } catch (IOException iox) { + LOGGER.warn("Could not close socket.", iox); + } + socket_ = null; + } + } + +} \ No newline at end of file From 370b64115525019bf63b860f0eea4e7a532d1a3c Mon Sep 17 00:00:00 2001 From: T Jake Luciani Date: Fri, 30 Sep 2011 18:20:28 +0000 Subject: [PATCH 11/15] Performance improvement for ByteBufferUtil.compareUnsigned Patch by tjake; reviewed by jbellis for CASSANDRA-3286 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0@1177766 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + src/java/org/apache/cassandra/utils/ByteBufferUtil.java | 6 ++++++ src/java/org/apache/cassandra/utils/FBUtilities.java | 6 +++--- 3 files changed, 10 insertions(+), 3 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 219bcb60d7..63074d3226 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,6 +1,7 @@ 1.0.1 * describe_ring should include datacenter/topology information (CASSANDRA-2882) * Thrift sockets are not properly buffered (CASSANDRA-3261) + * performance improvement for bytebufferutil compare function (CASSANDRA-3286) 1.0.0-final diff --git a/src/java/org/apache/cassandra/utils/ByteBufferUtil.java b/src/java/org/apache/cassandra/utils/ByteBufferUtil.java index 220fea0fc4..86c00c1879 100644 --- a/src/java/org/apache/cassandra/utils/ByteBufferUtil.java +++ b/src/java/org/apache/cassandra/utils/ByteBufferUtil.java @@ -83,6 +83,12 @@ public class ByteBufferUtil assert o1 != null; assert o2 != null; + if (o1.hasArray() && o2.hasArray()) + { + return FBUtilities.compareUnsigned(o1.array(), o2.array(), o1.position() + o1.arrayOffset(), + o2.position() + o2.arrayOffset(), o1.remaining(), o2.remaining()); + } + int minLength = Math.min(o1.remaining(), o2.remaining()); for (int x = 0, i = o1.position(), j = o2.position(); x < minLength; x++, i++, j++) { diff --git a/src/java/org/apache/cassandra/utils/FBUtilities.java b/src/java/org/apache/cassandra/utils/FBUtilities.java index 4a8aa71a6f..b1468863c0 100644 --- a/src/java/org/apache/cassandra/utils/FBUtilities.java +++ b/src/java/org/apache/cassandra/utils/FBUtilities.java @@ -348,7 +348,7 @@ public class FBUtilities } if (bytes2 == null) return 1; - int minLength = Math.min(len1 - offset1, len2 - offset2); + int minLength = Math.min(len1, len2); for (int x = 0, i = offset1, j = offset2; x < minLength; x++, i++, j++) { if (bytes1[i] == bytes2[j]) @@ -356,8 +356,8 @@ public class FBUtilities // compare non-equal bytes as unsigned return (bytes1[i] & 0xFF) < (bytes2[j] & 0xFF) ? -1 : 1; } - if ((len1 - offset1) == (len2 - offset2)) return 0; - else return ((len1 - offset1) < (len2 - offset2)) ? -1 : 1; + if (len1 == len2) return 0; + else return (len1 < len2) ? -1 : 1; } /** From 01eb6fec979c6b5a1b3101dfdb91b8a1358090e4 Mon Sep 17 00:00:00 2001 From: Pavel Yaskevich Date: Fri, 30 Sep 2011 20:20:23 +0000 Subject: [PATCH 12/15] CLI documentation change for ColumnFamily `compression_options` patch by Pavel Yaskevich; reviewed by Brandon Williams for CASSANDRA-3282 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177813 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 +- .../org/apache/cassandra/cli/CliHelp.yaml | 46 +++++++++++-------- 2 files changed, 28 insertions(+), 20 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 3e4ef70611..723c144c34 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -11,7 +11,7 @@ * fix counter entry in jdbc TypesMap (CASSANDRA-3268) * fix full queue scenario for ParallelCompactionIterator (CASSANDRA-3270) * fix bootstrap process (CASSANDRA-3285) - + * CLI documentation change for ColumnFamily `compression_options` (CASSANDRA-3282) 1.0.0-rc1 * Update CQL to generate microsecond timestamps by default (CASSANDRA-3227) diff --git a/src/resources/org/apache/cassandra/cli/CliHelp.yaml b/src/resources/org/apache/cassandra/cli/CliHelp.yaml index 981ebfbe65..8e1979cda3 100644 --- a/src/resources/org/apache/cassandra/cli/CliHelp.yaml +++ b/src/resources/org/apache/cassandra/cli/CliHelp.yaml @@ -563,25 +563,20 @@ commands: more rows in a given memory footprint. And storing the cache off-heap means you can use smaller heap sizes, reducing the impact of GC pauses. - - compression: Use compression for SSTable data files. - - Supported values are: - - null: to disable compression - - SnappyCompressor: compression based on the Snappy algorithm - - DeflateCompressor: compression based on the deflate algorithm - (through Java native support) - - It is also valid to specify the fully-qualified class name to a class - that implements org.apache.cassandra.io.ICompressor. - - compression_options: Options related to compression. - Options have the form [{key:value}]. The main recognized option are: - - sstable_compression: the algorithm to use to compress sstables for - this column family. If none is provided, compression will not be - enabled. Supported values are SnappyCompressor, DeflateCompressor or - any custom compressor. - - chunk_length_kb: specify the size of the chunk used by sstable - compression (default to 64, must be a power of 2). + Options have the form {key:value}. + The main recognized options are: + - sstable_compression: the algorithm to use to compress sstables for + this column family. If none is provided, compression will not be + enabled. Supported values are SnappyCompressor, DeflateCompressor or + any custom compressor. It is also valid to specify the fully-qualified + class name to a class that implements org.apache.cassandra.io.ICompressor. + + - chunk_length_kb: specify the size of the chunk used by sstable + compression (default to 64, must be a power of 2). + + To disable compression just set compression_options to null like this + `compression_options = null`. Examples: create column family Super4 @@ -836,7 +831,20 @@ commands: memory footprint. And storing the cache off-heap means you can use smaller heap sizes, reducing the impact of GC pauses. - - compression: Use compression for SSTable data files. Accepts the values true and false. + - compression_options: Options related to compression. + Options have the form {key:value}. + The main recognized options are: + - sstable_compression: the algorithm to use to compress sstables for + this column family. If none is provided, compression will not be + enabled. Supported values are SnappyCompressor, DeflateCompressor or + any custom compressor. It is also valid to specify the fully-qualified + class name to a class that implements org.apache.cassandra.io.ICompressor. + + - chunk_length_kb: specify the size of the chunk used by sstable + compression (default to 64, must be a power of 2). + + To disable compression just set compression_options to null like this + `compression_options = null`. Examples: update column family Super4 From 91c4cba8a58f229f3d51c33ce9a885d55758dedc Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 30 Sep 2011 20:45:44 +0000 Subject: [PATCH 13/15] remove obsolete memtable options from cli help patch by jbellis for CASSANDRA-3284 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177825 13f79535-47bb-0310-9956-ffa450edef68 --- .../org/apache/cassandra/cli/CliHelp.yaml | 16 ++-------------- 1 file changed, 2 insertions(+), 14 deletions(-) diff --git a/src/resources/org/apache/cassandra/cli/CliHelp.yaml b/src/resources/org/apache/cassandra/cli/CliHelp.yaml index 8e1979cda3..7414955743 100644 --- a/src/resources/org/apache/cassandra/cli/CliHelp.yaml +++ b/src/resources/org/apache/cassandra/cli/CliHelp.yaml @@ -472,14 +472,8 @@ commands: terms of I/O for the key cache. Row cache saving is much more expensive and has limited use. - - memtable_operations: Number of operations in millions before the memtable - is flushed. Default is memtable_throughput / 64 * 0.3 - - - memtable_throughput: Maximum size in MB to let a memtable get to before - it is flushed. Default is to use 1/16 the JVM heap size. - - read_repair_chance: Probability (0.0-1.0) with which to perform read - repairs for any read operation. Default is 1.0 to enable read repair. + repairs for any read operation. Default is 0.1. Note that disabling read repair entirely means that the dynamic snitch will not have any latency information from all the replicas to recognize @@ -739,14 +733,8 @@ commands: terms of I/O for the key cache. Row cache saving is much more expensive and has limited use. - - memtable_operations: Number of operations in millions before the memtable - is flushed. Default is memtable_throughput / 64 * 0.3 - - - memtable_throughput: Maximum size in MB to let a memtable get to before - it is flushed. Default is to use 1/16 the JVM heap size. - - read_repair_chance: Probability (0.0-1.0) with which to perform read - repairs for any read operation. Default is 1.0 to enable read repair. + repairs for any read operation. Default is 0.1. Note that disabling read repair entirely means that the dynamic snitch will not have any latency information from all the replicas to recognize From 160cb28d97b3d14094e60f53cb760f3f145f72e3 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Sat, 1 Oct 2011 02:46:51 +0000 Subject: [PATCH 14/15] ignore any CF ids sent by client for adding CF/KS patch by jbellis and Nate McCall for CASSANDRA-3288 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177887 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 4 ++++ src/java/org/apache/cassandra/thrift/CassandraServer.java | 2 ++ 2 files changed, 6 insertions(+) diff --git a/CHANGES.txt b/CHANGES.txt index 723c144c34..7389532bcb 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,3 +1,7 @@ +1.0.0-final + * ignore any CF ids sent by client for adding CF/KS (CASSANDRA-3288) + + 1.0.0-rc2 * Log a meaningful warning when a node receives a message for a repair session that doesn't exist anymore (CASSANDRA-3256) diff --git a/src/java/org/apache/cassandra/thrift/CassandraServer.java b/src/java/org/apache/cassandra/thrift/CassandraServer.java index 055bb2ff96..58984c5225 100644 --- a/src/java/org/apache/cassandra/thrift/CassandraServer.java +++ b/src/java/org/apache/cassandra/thrift/CassandraServer.java @@ -892,6 +892,7 @@ public class CassandraServer implements Cassandra.Iface try { + cf_def.unsetId(); // explicitly ignore any id set by client (Hector likes to set zero) applyMigrationOnStage(new AddColumnFamily(CFMetaData.fromThrift(cf_def))); return Schema.instance.getVersion().toString(); } @@ -957,6 +958,7 @@ public class CassandraServer implements Cassandra.Iface Collection cfDefs = new ArrayList(ks_def.cf_defs.size()); for (CfDef cf_def : ks_def.cf_defs) { + cf_def.unsetId(); // explicitly ignore any id set by client (same as system_add_column_family) CFMetaData.addDefaultIndexNames(cf_def); ThriftValidation.validateCfDef(cf_def, null); cfDefs.add(CFMetaData.fromThrift(cf_def)); From 2dbde4357ba52b73b4245feb12061718b24839d8 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Sat, 1 Oct 2011 05:41:09 +0000 Subject: [PATCH 15/15] remove obsolete hints on first startup patch by jbellis; reviewed by Patricio Echague for CASSANDRA-3291 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1177922 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../org/apache/cassandra/db/SystemTable.java | 34 ++++++++----------- 2 files changed, 16 insertions(+), 19 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 7389532bcb..b1cd3ed679 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,6 @@ 1.0.0-final * ignore any CF ids sent by client for adding CF/KS (CASSANDRA-3288) + * remove obsolete hints on first startup (CASSANDRA-3291) 1.0.0-rc2 diff --git a/src/java/org/apache/cassandra/db/SystemTable.java b/src/java/org/apache/cassandra/db/SystemTable.java index 70d52048d1..8ee22f2ea7 100644 --- a/src/java/org/apache/cassandra/db/SystemTable.java +++ b/src/java/org/apache/cassandra/db/SystemTable.java @@ -18,8 +18,6 @@ package org.apache.cassandra.db; -import java.io.File; -import java.io.FilenameFilter; import java.io.IOError; import java.io.IOException; import java.net.InetAddress; @@ -30,7 +28,6 @@ import java.util.List; import java.util.ArrayList; import java.util.SortedSet; import java.util.TreeSet; -import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.ExecutionException; import org.slf4j.Logger; @@ -74,25 +71,24 @@ public class SystemTable /* if hints become incompatible across versions of cassandra, that logic (and associated purging) is managed here. */ public static void purgeIncompatibleHints() throws IOException { - // 0.6->0.7 - final ByteBuffer hintsPurged6to7 = ByteBufferUtil.bytes("Hints purged as part of upgrading from 0.6.x to 0.7"); + ByteBuffer upgradeMarker = ByteBufferUtil.bytes("Pre-1.0 hints purged"); Table table = Table.open(Table.SYSTEM_TABLE); - QueryFilter dotSeven = QueryFilter.getNamesFilter(decorate(COOKIE_KEY), new QueryPath(STATUS_CF), hintsPurged6to7); - ColumnFamily cf = table.getColumnFamilyStore(STATUS_CF).getColumnFamily(dotSeven); - if (cf == null) + QueryFilter filter = QueryFilter.getNamesFilter(decorate(COOKIE_KEY), new QueryPath(STATUS_CF), upgradeMarker); + ColumnFamily cf = table.getColumnFamilyStore(STATUS_CF).getColumnFamily(filter); + if (cf != null) + return; + + // marker not found. Snapshot + remove hints and add the marker + ColumnFamilyStore hintsCfs = Table.open(Table.SYSTEM_TABLE).getColumnFamilyStore(HintedHandOffManager.HINTS_CF); + if (hintsCfs.getSSTables().size() > 0) { - // 0.7+ marker not found. Remove hints and add the marker. - ColumnFamilyStore hintsCfs = Table.open(Table.SYSTEM_TABLE).getColumnFamilyStore(HintedHandOffManager.HINTS_CF); - if (hintsCfs.getSSTables().size() > 0) - { - logger.info("Possible 0.6-format hints found. Snapshotting as 'old-hints' and purging"); - hintsCfs.snapshot("old-hints"); - hintsCfs.removeAllSSTables(); - } - RowMutation rm = new RowMutation(Table.SYSTEM_TABLE, COOKIE_KEY); - rm.add(new QueryPath(STATUS_CF, null, hintsPurged6to7), ByteBufferUtil.bytes("oh yes, it they were purged."), System.currentTimeMillis()); - rm.apply(); + logger.info("Possible old-format hints found. Snapshotting as 'old-hints' and purging"); + hintsCfs.snapshot("old-hints"); + hintsCfs.removeAllSSTables(); } + RowMutation rm = new RowMutation(Table.SYSTEM_TABLE, COOKIE_KEY); + rm.add(new QueryPath(STATUS_CF, null, upgradeMarker), ByteBufferUtil.bytes("oh yes, they were purged"), System.currentTimeMillis()); + rm.apply(); } /**