From 2f5f0c2c12a2d9c91cad79655b80382a2bd4c132 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 2 Dec 2010 01:48:18 +0000 Subject: [PATCH 01/22] fix consistencylevel calculations forNetworkTopologyStrategy patch by jbellis; reviewed by Jon Hermes for CASSANDRA-1804 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041250 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 ++ .../cassandra/config/DatabaseDescriptor.java | 5 ----- .../locator/AbstractReplicationStrategy.java | 5 ++--- .../service/QuorumResponseHandler.java | 5 +++-- .../cassandra/service/StorageService.java | 2 +- .../service/WriteResponseHandler.java | 5 +++-- .../cassandra/dht/BootStrapperTest.java | 3 ++- .../OldNetworkTopologyStrategyTest.java | 22 ++++++------------- .../cassandra/locator/SimpleStrategyTest.java | 4 ++-- .../service/AntiEntropyServiceTest.java | 4 ++-- .../apache/cassandra/service/MoveTest.java | 6 ++--- 11 files changed, 27 insertions(+), 36 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 376b556354..7342354997 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -23,6 +23,8 @@ dev * close file handle used for post-flush truncate (CASSANDRA-1790) * various code cleanup (CASSANDRA-1793, -1794, -1795) * fix range queries against wrapped range (CASSANDRA-1781) + * fix consistencylevel calculations for NetworkTopologyStrategy + (CASSANDRA-1804) 0.7.0-rc1 diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 0cae7d73e9..ab0415d112 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -815,11 +815,6 @@ public class DatabaseDescriptor return conf.rpc_port; } - public static int getReplicationFactor(String table) - { - return tables.get(table).replicationFactor; - } - public static long getRpcTimeout() { return conf.rpc_timeout_in_ms; diff --git a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java index 9ec15c114b..1f997cbafc 100644 --- a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java +++ b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java @@ -125,10 +125,9 @@ public abstract class AbstractReplicationStrategy return WriteResponseHandler.create(writeEndpoints, hintedEndpoints, consistencyLevel, table); } - // instance method so test subclasses can override it - int getReplicationFactor() + public int getReplicationFactor() { - return DatabaseDescriptor.getReplicationFactor(table); + return DatabaseDescriptor.getTableDefinition(table).replicationFactor; } /** diff --git a/src/java/org/apache/cassandra/service/QuorumResponseHandler.java b/src/java/org/apache/cassandra/service/QuorumResponseHandler.java index c21aac2e42..a703e05718 100644 --- a/src/java/org/apache/cassandra/service/QuorumResponseHandler.java +++ b/src/java/org/apache/cassandra/service/QuorumResponseHandler.java @@ -25,6 +25,7 @@ import java.util.concurrent.TimeoutException; import java.io.IOException; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.db.Table; import org.apache.cassandra.net.IAsyncCallback; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -110,9 +111,9 @@ public class QuorumResponseHandler implements IAsyncCallback case ANY: return 1; case QUORUM: - return (DatabaseDescriptor.getReplicationFactor(table) / 2) + 1; + return (Table.open(table).getReplicationStrategy().getReplicationFactor() / 2) + 1; case ALL: - return DatabaseDescriptor.getReplicationFactor(table); + return Table.open(table).getReplicationStrategy().getReplicationFactor(); default: throw new UnsupportedOperationException("invalid consistency level: " + table.toString()); } diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 7471aaf55d..2bb83f9838 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -1738,7 +1738,7 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe for (String table : DatabaseDescriptor.getNonSystemTables()) { // if the replication factor is 1 the data is lost so we shouldn't wait for confirmation - if (DatabaseDescriptor.getReplicationFactor(table) == 1) + if (Table.open(table).getReplicationStrategy().getReplicationFactor() == 1) continue; // get all ranges that change ownership (that is, a node needs diff --git a/src/java/org/apache/cassandra/service/WriteResponseHandler.java b/src/java/org/apache/cassandra/service/WriteResponseHandler.java index 3ea15c97f4..95e5e3f362 100644 --- a/src/java/org/apache/cassandra/service/WriteResponseHandler.java +++ b/src/java/org/apache/cassandra/service/WriteResponseHandler.java @@ -26,6 +26,7 @@ import java.util.concurrent.atomic.AtomicInteger; import com.google.common.collect.ImmutableMultimap; import com.google.common.collect.Multimap; import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.db.Table; import org.apache.cassandra.net.Message; import org.apache.cassandra.thrift.ConsistencyLevel; import org.apache.cassandra.thrift.UnavailableException; @@ -93,9 +94,9 @@ public class WriteResponseHandler extends AbstractWriteResponseHandler } // at most one node per range can bootstrap at a time, and these will be added to the write until // bootstrap finishes (at which point we no longer need to write to the old ones). - assert 1 <= blockFor && blockFor <= 2 * DatabaseDescriptor.getReplicationFactor(table) + assert 1 <= blockFor && blockFor <= 2 * Table.open(table).getReplicationStrategy().getReplicationFactor() : String.format("invalid response count %d for replication factor %d", - blockFor, DatabaseDescriptor.getReplicationFactor(table)); + blockFor, Table.open(table).getReplicationStrategy().getReplicationFactor()); return blockFor; } diff --git a/test/unit/org/apache/cassandra/dht/BootStrapperTest.java b/test/unit/org/apache/cassandra/dht/BootStrapperTest.java index 7fc460203a..d42e87c270 100644 --- a/test/unit/org/apache/cassandra/dht/BootStrapperTest.java +++ b/test/unit/org/apache/cassandra/dht/BootStrapperTest.java @@ -33,6 +33,7 @@ import org.junit.Test; import com.google.common.collect.Multimap; +import org.apache.cassandra.db.Table; import org.apache.cassandra.gms.ApplicationState; import org.apache.cassandra.gms.IFailureDetectionEventListener; import org.apache.cassandra.gms.IFailureDetector; @@ -146,7 +147,7 @@ public class BootStrapperTest extends CleanupHelper final int[] clusterSizes = new int[] { 1, 3, 5, 10, 100}; for (String table : DatabaseDescriptor.getNonSystemTables()) { - int replicationFactor = DatabaseDescriptor.getReplicationFactor(table); + int replicationFactor = Table.open(table).getReplicationStrategy().getReplicationFactor(); for (int clusterSize : clusterSizes) if (clusterSize >= replicationFactor) testSourceTargetComputation(table, clusterSize, replicationFactor); diff --git a/test/unit/org/apache/cassandra/locator/OldNetworkTopologyStrategyTest.java b/test/unit/org/apache/cassandra/locator/OldNetworkTopologyStrategyTest.java index c42d7c3e2c..288491d3a8 100644 --- a/test/unit/org/apache/cassandra/locator/OldNetworkTopologyStrategyTest.java +++ b/test/unit/org/apache/cassandra/locator/OldNetworkTopologyStrategyTest.java @@ -30,11 +30,12 @@ import org.junit.Before; import org.junit.Test; import static org.junit.Assert.assertEquals; -import org.apache.cassandra.config.DatabaseDescriptor; + +import org.apache.cassandra.SchemaLoader; import org.apache.cassandra.dht.BigIntegerToken; import org.apache.cassandra.dht.Token; -public class OldNetworkTopologyStrategyTest +public class OldNetworkTopologyStrategyTest extends SchemaLoader { private List endpointTokens; private List keyTokens; @@ -71,7 +72,7 @@ public class OldNetworkTopologyStrategyTest expectedResults.put("25", buildResult("254.0.0.4", "254.0.0.1", "254.0.0.2")); expectedResults.put("35", buildResult("254.0.0.1", "254.0.0.2", "254.0.0.3")); - runTestForReplicatedTables(strategy); + testGetEndpoints(strategy, keyTokens.toArray(new Token[0])); } /** @@ -96,7 +97,7 @@ public class OldNetworkTopologyStrategyTest expectedResults.put("25", buildResult("254.0.0.4", "254.1.0.3", "254.0.0.1")); expectedResults.put("35", buildResult("254.0.0.1", "254.1.0.3", "254.0.0.2")); - runTestForReplicatedTables(strategy); + testGetEndpoints(strategy, keyTokens.toArray(new Token[0])); } /** @@ -122,16 +123,7 @@ public class OldNetworkTopologyStrategyTest expectedResults.put("25", buildResult("254.1.0.4", "254.0.0.1", "254.0.0.2")); expectedResults.put("35", buildResult("254.0.0.1", "254.0.1.3", "254.1.0.4")); - runTestForReplicatedTables(strategy); - } - - private void runTestForReplicatedTables(AbstractReplicationStrategy strategy) throws UnknownHostException - { - for (String table : DatabaseDescriptor.getNonSystemTables()) - { - if (DatabaseDescriptor.getReplicationFactor(table) == 3) - testGetEndpoints(strategy, keyTokens.toArray(new Token[0]), table); - } + testGetEndpoints(strategy, keyTokens.toArray(new Token[0])); } private ArrayList buildResult(String... addresses) throws UnknownHostException @@ -156,7 +148,7 @@ public class OldNetworkTopologyStrategyTest tmd.updateNormalToken(endpointToken, ep); } - private void testGetEndpoints(AbstractReplicationStrategy strategy, Token[] keyTokens, String table) throws UnknownHostException + private void testGetEndpoints(AbstractReplicationStrategy strategy, Token[] keyTokens) throws UnknownHostException { for (Token keyToken : keyTokens) { diff --git a/test/unit/org/apache/cassandra/locator/SimpleStrategyTest.java b/test/unit/org/apache/cassandra/locator/SimpleStrategyTest.java index f53680c783..56d7089ac4 100644 --- a/test/unit/org/apache/cassandra/locator/SimpleStrategyTest.java +++ b/test/unit/org/apache/cassandra/locator/SimpleStrategyTest.java @@ -94,7 +94,7 @@ public class SimpleStrategyTest extends CleanupHelper for (int i = 0; i < keyTokens.length; i++) { List endpoints = strategy.getNaturalEndpoints(keyTokens[i]); - assertEquals(DatabaseDescriptor.getReplicationFactor(table), endpoints.size()); + assertEquals(strategy.getReplicationFactor(), endpoints.size()); List correctEndpoints = new ArrayList(); for (int j = 0; j < endpoints.size(); j++) correctEndpoints.add(hosts.get((i + j + 1) % hosts.size())); @@ -140,7 +140,7 @@ public class SimpleStrategyTest extends CleanupHelper StorageService.calculatePendingRanges(strategy, table); - int replicationFactor = DatabaseDescriptor.getReplicationFactor(table); + int replicationFactor = strategy.getReplicationFactor(); for (int i = 0; i < keyTokens.length; i++) { diff --git a/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java b/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java index ccbdcc0157..7402233469 100644 --- a/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java +++ b/test/unit/org/apache/cassandra/service/AntiEntropyServiceTest.java @@ -171,7 +171,7 @@ public class AntiEntropyServiceTest extends CleanupHelper public void testGetNeighborsPlusOne() throws Throwable { // generate rf+1 nodes, and ensure that all nodes are returned - Set expected = addTokens(1 + DatabaseDescriptor.getReplicationFactor(tablename)); + Set expected = addTokens(1 + Table.open(tablename).getReplicationStrategy().getReplicationFactor()); expected.remove(FBUtilities.getLocalAddress()); assertEquals(expected, AntiEntropyService.getNeighbors(tablename)); } @@ -182,7 +182,7 @@ public class AntiEntropyServiceTest extends CleanupHelper TokenMetadata tmd = StorageService.instance.getTokenMetadata(); // generate rf*2 nodes, and ensure that only neighbors specified by the ARS are returned - addTokens(2 * DatabaseDescriptor.getReplicationFactor(tablename)); + addTokens(2 * Table.open(tablename).getReplicationStrategy().getReplicationFactor()); AbstractReplicationStrategy ars = Table.open(tablename).getReplicationStrategy(); Set expected = new HashSet(); for (Range replicaRange : ars.getAddressRanges().get(FBUtilities.getLocalAddress())) diff --git a/test/unit/org/apache/cassandra/service/MoveTest.java b/test/unit/org/apache/cassandra/service/MoveTest.java index ba75a22d93..90b80ef7c3 100644 --- a/test/unit/org/apache/cassandra/service/MoveTest.java +++ b/test/unit/org/apache/cassandra/service/MoveTest.java @@ -92,7 +92,7 @@ public class MoveTest extends CleanupHelper strategy = getStrategy(table, tmd); for (Token token : keyTokens) { - int replicationFactor = DatabaseDescriptor.getReplicationFactor(table); + int replicationFactor = strategy.getReplicationFactor(); HashSet actual = new HashSet(tmd.getWriteEndpoints(token, table, strategy.calculateNaturalEndpoints(token, tmd))); HashSet expected = new HashSet(); @@ -217,7 +217,7 @@ public class MoveTest extends CleanupHelper } // just to be sure that things still work according to the old tests, run them: - if (DatabaseDescriptor.getReplicationFactor(table) != 3) + if (strategy.getReplicationFactor() != 3) continue; // tokens 5, 15 and 25 should go three nodes for (int i=0; i<3; ++i) @@ -334,7 +334,7 @@ public class MoveTest extends CleanupHelper assertTrue(expectedEndpoints.get(table).get(keyTokens.get(i)).containsAll(endpoints)); } - if (DatabaseDescriptor.getReplicationFactor(table) != 3) + if (strategy.getReplicationFactor() != 3) continue; // leave this stuff in to guarantee the old tests work the way they were supposed to. // tokens 5, 15 and 25 should go three nodes From af5e0b7dd449d3cfdd941b54896e6fe7fe623f4a Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 2 Dec 2010 17:39:36 +0000 Subject: [PATCH 04/22] r/m out-of-date contrib/maven for CASSANDRA-1805 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041486 13f79535-47bb-0310-9956-ffa450edef68 --- contrib/maven/pom.xml | 551 ------------------ .../org/apache/cassandra/db/RowMutation.java | 1 + .../cassandra/db/RowMutationVerbHandler.java | 50 +- .../cassandra/service/StorageProxy.java | 178 ++++-- 4 files changed, 176 insertions(+), 604 deletions(-) delete mode 100644 contrib/maven/pom.xml diff --git a/contrib/maven/pom.xml b/contrib/maven/pom.xml deleted file mode 100644 index 5d98eb76a8..0000000000 --- a/contrib/maven/pom.xml +++ /dev/null @@ -1,551 +0,0 @@ - - - - - org.apache - apache - 6 - - 4.0.0 - - org.apache.cassandra - cassandra - 0.6-SNAPSHOT - jar - Cassandra - 2009 - - - 2.0.9 - - - http://incubator.apache.org/cassandra - - - - cassandra-user - cassandra-user-subscribe@incubator.apache.org - cassandra-user-unsubscribe@incubator.apache.org - cassandra-user@incubator.apache.org - http://mail-archives.apache.org/mod_mbox/incubator-cassandra-user/ - - - cassandra-dev - cassandra-dev-subscribe@incubator.apache.org - cassandra-dev-unsubscribe@incubator.apache.org - cassandra-dev@incubator.apache.org - http://mail-archives.apache.org/mod_mbox/incubator-cassandra-dev/ - - - cassandra-commits - cassandra-commits-subscribe@incubator.apache.org - cassandra-commits-unsubscribe@incubator.apache.org - cassandra-commits@incubator.apache.org - http://mail-archives.apache.org/mod_mbox/incubator-cassandra-commits/ - - - - - - - - - - - JIRA - https://issues.apache.org/jira/browse/CASSANDRA - - - - - hudson - http://hudson.zones.apache.org/hudson - - - - mail - true - true - false - false -
cassandra-commits@incubator.apache.org
-
-
-
- - - - cassandra-website - scp://people.apache.org/x1/www/incubator.apache.org/cassandra/maven/${pom.version} - - - - - scm:svn:http://svn.apache.org/repos/asf/incubator/cassandra/trunk - scm:svn:https://svn.apache.org/repos/asf/incubator/cassandra/trunk - http://svn.apache.org/viewvc/incubator/cassandra/trunk/ - - - - - - commons-collections - commons-collections - 3.2.1 - - - commons-cli - commons-cli - 1.1 - - - commons-lang - commons-lang - 2.4 - - - jline - jline - 0.9.94 - - - log4j - log4j - 1.2.15 - - - - javax.jms - jms - - - com.sun.jmx - jmxri - - - com.sun.jdmk - jmxtools - - - - - org.slf4j - slf4j-api - 1.5.8 - - - org.slf4j - slf4j-log4j12 - 1.5.8 - - - org.antlr - antlr-runtime - 3.1.3 - - - com.google.collections - google-collections - 1.0-rc1 - - - - - - high-scale-lib - high-scale-lib - UNKNOWN - system - ${basedir}/lib/high-scale-lib.jar - - - flexjson - flexjson - 1.7 - system - ${basedir}/lib/flexjson-1.7.jar - - - libthrift - libthrift - UNKNOWN - system - ${basedir}/lib/libthrift-r820831.jar - - - jsonsimple - jsonsimple - UNKNOWN - system - ${basedir}/lib/json_simple-1.1.jar - - - com.reardencommerce - clhm - UNKNOWN - system - ${basedir}/lib/clhm-production.jar - - - - - junit - junit - 4.6 - test - - - - - - ${basedir}/src/java - ${basedir}/test/unit - build/classes - - - - ${basedir}/test/conf - - **/* - - - - ${basedir}/test/resources - - *.json - - - - - - - - - org.antlr - antlr3-maven-plugin - 3.1.3-1 - - - process-sources - - antlr - - - ${basedir}/src/java - - - - - - - - org.codehaus.mojo - build-helper-maven-plugin - 1.3 - - - add-source - generate-sources - - add-source - - - - ${basedir}/interface/gen-java - - - - - - - - - org.apache.maven.plugins - maven-compiler-plugin - - 1.6 - 1.6 - true - true - true - true - - - - - - org.apache.maven.plugins - maven-surefire-plugin - - - - storage-config - ${basedir}/test/conf - - - always - - **/TestRingCache.java - - - - - - - org.codehaus.mojo - cobertura-maven-plugin - 2.0 - - - - - - - org.apache.maven.plugins - maven-release-plugin - 2.0-beta-9 - - true - false - clean install - deploy - -Papache-release - - - - org.codehaus.mojo - ianal-maven-plugin - 1.0-alpha-1 - - - org.codehaus.mojo - rat-maven-plugin - 1.0-alpha-3 - - false - - - - org.apache.maven.plugins - maven-enforcer-plugin - - - validate - - enforce - - - - - [2.0.9,) - - - - - - - - org.codehaus.mojo - ianal-maven-plugin - - - - verify-legal-files - - - true - - - - - - - - - - - org.apache.maven.plugins - maven-jxr-plugin - - - org.apache.maven.plugins - maven-surefire-report-plugin - - - org.apache.maven.plugins - maven-pmd-plugin - - - org.codehaus.mojo - taglist-maven-plugin - - - org.apache.maven.plugins - maven-javadoc-plugin - - - http://java.sun.com/j2se/1.6.0/docs/api/ - http://logging.apache.org/log4j/docs/api/ - - - true - 900m - 1.6 - - - - org.codehaus.mojo - cobertura-maven-plugin - 2.2 - - - html - xml - - - - - - - - - - - thrift - - process-sources - - - - maven-antrun-plugin - - - process-sources - - - - - - - - run - - - - - - - - - - - apache-release - - - - - - org.apache.maven.plugins - maven-gpg-plugin - - ${gpg.passphrase} - - - - - sign - - - - - - - - true - org.apache.maven.plugins - maven-deploy-plugin - - true - - - - org.apache.maven.plugins - maven-source-plugin - - - attach-sources - - jar - - - - - - org.apache.maven.plugins - maven-javadoc-plugin - - ${project.build.sourceEncoding} - - - - attach-javadocs - - jar - - - - - - - - org.apache.maven.plugins - maven-assembly-plugin - - - - single - - package - - true - - - source-release - - - - - - - - org.apache.geronimo.genesis - apache-source-release-assembly-descriptor - 2.0 - - - - - - - - - - -
diff --git a/src/java/org/apache/cassandra/db/RowMutation.java b/src/java/org/apache/cassandra/db/RowMutation.java index 5217d83ac0..60e0a0177b 100644 --- a/src/java/org/apache/cassandra/db/RowMutation.java +++ b/src/java/org/apache/cassandra/db/RowMutation.java @@ -47,6 +47,7 @@ public class RowMutation { private static ICompactSerializer serializer_; public static final String HINT = "HINT"; + public static final String FORWARD_HEADER = "FORWARD"; static { diff --git a/src/java/org/apache/cassandra/db/RowMutationVerbHandler.java b/src/java/org/apache/cassandra/db/RowMutationVerbHandler.java index bd6d8f29d2..92b89b4b5d 100644 --- a/src/java/org/apache/cassandra/db/RowMutationVerbHandler.java +++ b/src/java/org/apache/cassandra/db/RowMutationVerbHandler.java @@ -18,25 +18,23 @@ package org.apache.cassandra.db; -import java.io.*; - +import java.io.ByteArrayInputStream; +import java.io.DataInputStream; +import java.io.IOException; import java.net.InetAddress; +import java.net.UnknownHostException; import java.nio.ByteBuffer; import com.google.common.base.Charsets; - -import org.apache.cassandra.net.IVerbHandler; -import org.apache.cassandra.net.Message; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.net.*; +import org.apache.cassandra.net.IVerbHandler; +import org.apache.cassandra.net.Message; +import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; -import static com.google.common.base.Charsets.UTF_8; - public class RowMutationVerbHandler implements IVerbHandler { @@ -69,6 +67,11 @@ public class RowMutationVerbHandler implements IVerbHandler hintedMutation.apply(); } } + + // Check if there were any forwarding headers in this message + byte[] forwardBytes = message.getHeader(RowMutation.FORWARD_HEADER); + if (forwardBytes != null) + forwardToLocalNodes(message, forwardBytes); Table.open(rm.getTable()).apply(rm, bytes, true); @@ -82,5 +85,34 @@ public class RowMutationVerbHandler implements IVerbHandler { logger_.error("Error in row mutation", e); } + } + + private void forwardToLocalNodes(Message message, byte[] forwardBytes) throws UnknownHostException + { + // remove fwds from message to avoid infinite loop + message.setHeader(RowMutation.FORWARD_HEADER, null); + + int bytesPerInetAddress = FBUtilities.getLocalAddress().getAddress().length; + assert forwardBytes.length >= bytesPerInetAddress; + assert forwardBytes.length % bytesPerInetAddress == 0; + + int offset = 0; + byte[] addressBytes = new byte[bytesPerInetAddress]; + + // Send a message to each of the addresses on our Forward List + while (offset < forwardBytes.length) + { + System.arraycopy(forwardBytes, offset, addressBytes, 0, bytesPerInetAddress); + InetAddress address = InetAddress.getByAddress(addressBytes); + + if (logger_.isDebugEnabled()) + logger_.debug("Forwarding message to " + address); + + // Send the original message to the address specified by the FORWARD_HINT + // Let the response go back to the coordinator + MessagingService.instance.sendOneWay(message, message.getFrom()); + + offset += bytesPerInetAddress; + } } } diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 9b96607d55..cc652ae3ee 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -27,20 +27,23 @@ import java.util.concurrent.*; import javax.management.MBeanServer; import javax.management.ObjectName; +import com.google.common.collect.HashMultimap; import com.google.common.collect.Multimap; -import static com.google.common.base.Charsets.UTF_8; import org.apache.commons.lang.ArrayUtils; import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; - import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.concurrent.StageManager; import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.*; -import org.apache.cassandra.dht.*; +import org.apache.cassandra.db.filter.QueryFilter; +import org.apache.cassandra.dht.AbstractBounds; +import org.apache.cassandra.dht.Bounds; +import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.dht.Token; import org.apache.cassandra.gms.Gossiper; import org.apache.cassandra.locator.AbstractReplicationStrategy; import org.apache.cassandra.locator.TokenMetadata; @@ -53,7 +56,8 @@ import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.LatencyTracker; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.WrappedRunnable; -import org.apache.cassandra.db.filter.QueryFilter; + +import static com.google.common.base.Charsets.UTF_8; public class StorageProxy implements StorageProxyMBean { @@ -90,13 +94,14 @@ public class StorageProxy implements StorageProxyMBean * @param mutations the mutations to be applied across the replicas * @param consistency_level the consistency level for the operation */ - public static void mutate(List mutations, ConsistencyLevel consistency_level) throws UnavailableException, TimeoutException + public static void mutate(List mutations, ConsistencyLevel consistencyLevel) throws UnavailableException, TimeoutException { long startTime = System.nanoTime(); - ArrayList responseHandlers = new ArrayList(); + List responseHandlers = new ArrayList(); RowMutation mostRecentRowMutation = null; StorageService ss = StorageService.instance; + String localDataCenter = getDataCenter(FBUtilities.getLocalAddress()); try { @@ -110,58 +115,67 @@ public class StorageProxy implements StorageProxyMBean Collection writeEndpoints = ss.getTokenMetadata().getWriteEndpoints(StorageService.getPartitioner().getToken(rm.key()), table, naturalEndpoints); Multimap hintedEndpoints = rs.getHintedEndpoints(writeEndpoints); - // send out the writes, as in mutate() above, but this time with a callback that tracks responses - final IWriteResponseHandler responseHandler = rs.getWriteResponseHandler(writeEndpoints, hintedEndpoints, consistency_level); + final IWriteResponseHandler responseHandler = rs.getWriteResponseHandler(writeEndpoints, hintedEndpoints, consistencyLevel); + + // exit early if we can't fulfuill the CL at this time responseHandler.assureSufficientLiveNodes(); - + responseHandlers.add(responseHandler); - Message unhintedMessage = null; - for (Map.Entry> entry : hintedEndpoints.asMap().entrySet()) + + // Creates a Multimap that holds onto all the messages and addresses meant for a specific datacenter. + Multimap> dcMap = groupEndpointsByDataCenter(rm, hintedEndpoints, responseHandler); + + // Traverse all dataCenters where messages will be sent to. + for (Map.Entry>> entry : dcMap.asMap().entrySet()) { - InetAddress destination = entry.getKey(); - Collection targets = entry.getValue(); + String dataCenter = entry.getKey(); + + // Grab a set of all the messages bound for this dataCenter and create an iterator over this set. + Collection> messagesForDataCenter = entry.getValue(); + Iterator> iter = messagesForDataCenter.iterator(); + assert iter.hasNext(); - if (targets.size() == 1 && targets.iterator().next().equals(destination)) + // First endpoint in list is the destination for this group + Pair messageAndDestination = iter.next(); + + Message primaryMessage = messageAndDestination.left; + InetAddress target = messageAndDestination.right; + + // Add all the other destinations that are bound for the same dataCenter as a header in the primary message. + while (iter.hasNext()) { - // unhinted writes - if (destination.equals(FBUtilities.getLocalAddress())) + messageAndDestination = iter.next(); + assert messageAndDestination.left == primaryMessage; + + if (dataCenter.equals(localDataCenter)) { - insertLocalMessage(rm, responseHandler); + // direct write to local DC + assert primaryMessage.getHeader(RowMutation.FORWARD_HEADER) == null; + MessagingService.instance.sendOneWay(primaryMessage, target); } else { - // belongs on a different server. send it there. - if (unhintedMessage == null) - { - unhintedMessage = rm.makeRowMutationMessage(); - MessagingService.instance.addCallback(responseHandler, unhintedMessage.getMessageId()); - } - if (logger.isDebugEnabled()) - logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + unhintedMessage.getMessageId() + "@" + destination); - MessagingService.instance.sendOneWay(unhintedMessage, destination); + // group all nodes in this DC as forward headers on the primary message + ByteArrayOutputStream bos = new ByteArrayOutputStream(); + DataOutputStream dos = new DataOutputStream(bos); + + // append to older addresses + byte[] previousHints = primaryMessage.getHeader(RowMutation.FORWARD_HEADER); + if (previousHints != null) + dos.write(previousHints); + + dos.write(messageAndDestination.right.getAddress()); + primaryMessage.setHeader(RowMutation.FORWARD_HEADER, bos.toByteArray()); } - } - else - { - // hinted - Message hintedMessage = rm.makeRowMutationMessage(); - for (InetAddress target : targets) - { - if (!target.equals(destination)) - { - addHintHeader(hintedMessage, target); - if (logger.isDebugEnabled()) - logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + hintedMessage.getMessageId() + "@" + destination + " for " + target); - } - } - responseHandler.addHintCallback(hintedMessage, destination); - MessagingService.instance.sendOneWay(hintedMessage, destination); - } + } + + MessagingService.instance.sendOneWay(primaryMessage, target); } } + // wait for writes. throws timeoutexception if necessary for (IWriteResponseHandler responseHandler : responseHandlers) - { + { responseHandler.get(); } } @@ -178,6 +192,66 @@ public class StorageProxy implements StorageProxyMBean } } + + private static Multimap> groupEndpointsByDataCenter(RowMutation rm, Multimap endpoints, final IWriteResponseHandler responseHandler) throws IOException + { + + Set>> endpointSet = endpoints.asMap().entrySet(); + Multimap> dcMap = HashMultimap.create(endpointSet.size(), 10); + Message unhintedMessage = null; + + for (Map.Entry> entry : endpointSet) + { + InetAddress destination = entry.getKey(); + Collection targets = entry.getValue(); + + String dataCenter = getDataCenter(destination); + + if (targets.size() == 1 && targets.iterator().next().equals(destination)) + { + // unhinted writes + if (destination.equals(FBUtilities.getLocalAddress())) + { + insertLocalMessage(rm, responseHandler); + } + else + { + // belongs on a different server. + if (unhintedMessage == null) + { + unhintedMessage = rm.makeRowMutationMessage(); + MessagingService.instance.addCallback(responseHandler, unhintedMessage.getMessageId()); + } + + if (logger.isDebugEnabled()) + logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + unhintedMessage.getMessageId() + "@" + destination); + + dcMap.put(dataCenter, new Pair(unhintedMessage, destination)); + } + } + else + { + // hinted + Message hintedMessage = rm.makeRowMutationMessage(); + + for (InetAddress target : targets) + { + if (!target.equals(destination)) + { + addHintHeader(hintedMessage, target); + if (logger.isDebugEnabled()) + logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + hintedMessage.getMessageId() + "@" + destination + " for " + target); + } + } + + responseHandler.addHintCallback(hintedMessage, destination); + dcMap.put(dataCenter, new Pair(hintedMessage, destination)); + } + } + + return dcMap; + } + private static void addHintHeader(Message message, InetAddress target) throws IOException { @@ -192,6 +266,22 @@ public class StorageProxy implements StorageProxyMBean message.setHeader(RowMutation.HINT, bos.toByteArray()); } + private static String getDataCenter(InetAddress addr) + { + String dataCenter = null; + try + { + dataCenter = DatabaseDescriptor.getEndpointSnitch().getDatacenter(addr); + } + catch (UnsupportedOperationException e) + { + // SimpleSnitch throws this + dataCenter = "default"; + } + + return dataCenter; + } + private static void insertLocalMessage(final RowMutation rm, final IWriteResponseHandler responseHandler) { if (logger.isDebugEnabled()) From f0104cdce5ec423ea4b47f0b3bbd8a7b8c8440a6 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 2 Dec 2010 17:42:15 +0000 Subject: [PATCH 05/22] revert last git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041489 13f79535-47bb-0310-9956-ffa450edef68 --- contrib/maven/pom.xml | 551 ++++++++++++++++++ .../org/apache/cassandra/db/RowMutation.java | 1 - .../cassandra/db/RowMutationVerbHandler.java | 50 +- .../cassandra/service/StorageProxy.java | 178 ++---- 4 files changed, 604 insertions(+), 176 deletions(-) create mode 100644 contrib/maven/pom.xml diff --git a/contrib/maven/pom.xml b/contrib/maven/pom.xml new file mode 100644 index 0000000000..5d98eb76a8 --- /dev/null +++ b/contrib/maven/pom.xml @@ -0,0 +1,551 @@ + + + + + org.apache + apache + 6 + + 4.0.0 + + org.apache.cassandra + cassandra + 0.6-SNAPSHOT + jar + Cassandra + 2009 + + + 2.0.9 + + + http://incubator.apache.org/cassandra + + + + cassandra-user + cassandra-user-subscribe@incubator.apache.org + cassandra-user-unsubscribe@incubator.apache.org + cassandra-user@incubator.apache.org + http://mail-archives.apache.org/mod_mbox/incubator-cassandra-user/ + + + cassandra-dev + cassandra-dev-subscribe@incubator.apache.org + cassandra-dev-unsubscribe@incubator.apache.org + cassandra-dev@incubator.apache.org + http://mail-archives.apache.org/mod_mbox/incubator-cassandra-dev/ + + + cassandra-commits + cassandra-commits-subscribe@incubator.apache.org + cassandra-commits-unsubscribe@incubator.apache.org + cassandra-commits@incubator.apache.org + http://mail-archives.apache.org/mod_mbox/incubator-cassandra-commits/ + + + + + + + + + + + JIRA + https://issues.apache.org/jira/browse/CASSANDRA + + + + + hudson + http://hudson.zones.apache.org/hudson + + + + mail + true + true + false + false +
cassandra-commits@incubator.apache.org
+
+
+
+ + + + cassandra-website + scp://people.apache.org/x1/www/incubator.apache.org/cassandra/maven/${pom.version} + + + + + scm:svn:http://svn.apache.org/repos/asf/incubator/cassandra/trunk + scm:svn:https://svn.apache.org/repos/asf/incubator/cassandra/trunk + http://svn.apache.org/viewvc/incubator/cassandra/trunk/ + + + + + + commons-collections + commons-collections + 3.2.1 + + + commons-cli + commons-cli + 1.1 + + + commons-lang + commons-lang + 2.4 + + + jline + jline + 0.9.94 + + + log4j + log4j + 1.2.15 + + + + javax.jms + jms + + + com.sun.jmx + jmxri + + + com.sun.jdmk + jmxtools + + + + + org.slf4j + slf4j-api + 1.5.8 + + + org.slf4j + slf4j-log4j12 + 1.5.8 + + + org.antlr + antlr-runtime + 3.1.3 + + + com.google.collections + google-collections + 1.0-rc1 + + + + + + high-scale-lib + high-scale-lib + UNKNOWN + system + ${basedir}/lib/high-scale-lib.jar + + + flexjson + flexjson + 1.7 + system + ${basedir}/lib/flexjson-1.7.jar + + + libthrift + libthrift + UNKNOWN + system + ${basedir}/lib/libthrift-r820831.jar + + + jsonsimple + jsonsimple + UNKNOWN + system + ${basedir}/lib/json_simple-1.1.jar + + + com.reardencommerce + clhm + UNKNOWN + system + ${basedir}/lib/clhm-production.jar + + + + + junit + junit + 4.6 + test + + + + + + ${basedir}/src/java + ${basedir}/test/unit + build/classes + + + + ${basedir}/test/conf + + **/* + + + + ${basedir}/test/resources + + *.json + + + + + + + + + org.antlr + antlr3-maven-plugin + 3.1.3-1 + + + process-sources + + antlr + + + ${basedir}/src/java + + + + + + + + org.codehaus.mojo + build-helper-maven-plugin + 1.3 + + + add-source + generate-sources + + add-source + + + + ${basedir}/interface/gen-java + + + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + + 1.6 + 1.6 + true + true + true + true + + + + + + org.apache.maven.plugins + maven-surefire-plugin + + + + storage-config + ${basedir}/test/conf + + + always + + **/TestRingCache.java + + + + + + + org.codehaus.mojo + cobertura-maven-plugin + 2.0 + + + + + + + org.apache.maven.plugins + maven-release-plugin + 2.0-beta-9 + + true + false + clean install + deploy + -Papache-release + + + + org.codehaus.mojo + ianal-maven-plugin + 1.0-alpha-1 + + + org.codehaus.mojo + rat-maven-plugin + 1.0-alpha-3 + + false + + + + org.apache.maven.plugins + maven-enforcer-plugin + + + validate + + enforce + + + + + [2.0.9,) + + + + + + + + org.codehaus.mojo + ianal-maven-plugin + + + + verify-legal-files + + + true + + + + + + + + + + + org.apache.maven.plugins + maven-jxr-plugin + + + org.apache.maven.plugins + maven-surefire-report-plugin + + + org.apache.maven.plugins + maven-pmd-plugin + + + org.codehaus.mojo + taglist-maven-plugin + + + org.apache.maven.plugins + maven-javadoc-plugin + + + http://java.sun.com/j2se/1.6.0/docs/api/ + http://logging.apache.org/log4j/docs/api/ + + + true + 900m + 1.6 + + + + org.codehaus.mojo + cobertura-maven-plugin + 2.2 + + + html + xml + + + + + + + + + + + thrift + + process-sources + + + + maven-antrun-plugin + + + process-sources + + + + + + + + run + + + + + + + + + + + apache-release + + + + + + org.apache.maven.plugins + maven-gpg-plugin + + ${gpg.passphrase} + + + + + sign + + + + + + + + true + org.apache.maven.plugins + maven-deploy-plugin + + true + + + + org.apache.maven.plugins + maven-source-plugin + + + attach-sources + + jar + + + + + + org.apache.maven.plugins + maven-javadoc-plugin + + ${project.build.sourceEncoding} + + + + attach-javadocs + + jar + + + + + + + + org.apache.maven.plugins + maven-assembly-plugin + + + + single + + package + + true + + + source-release + + + + + + + + org.apache.geronimo.genesis + apache-source-release-assembly-descriptor + 2.0 + + + + + + + + + + +
diff --git a/src/java/org/apache/cassandra/db/RowMutation.java b/src/java/org/apache/cassandra/db/RowMutation.java index 60e0a0177b..5217d83ac0 100644 --- a/src/java/org/apache/cassandra/db/RowMutation.java +++ b/src/java/org/apache/cassandra/db/RowMutation.java @@ -47,7 +47,6 @@ public class RowMutation { private static ICompactSerializer serializer_; public static final String HINT = "HINT"; - public static final String FORWARD_HEADER = "FORWARD"; static { diff --git a/src/java/org/apache/cassandra/db/RowMutationVerbHandler.java b/src/java/org/apache/cassandra/db/RowMutationVerbHandler.java index 92b89b4b5d..bd6d8f29d2 100644 --- a/src/java/org/apache/cassandra/db/RowMutationVerbHandler.java +++ b/src/java/org/apache/cassandra/db/RowMutationVerbHandler.java @@ -18,23 +18,25 @@ package org.apache.cassandra.db; -import java.io.ByteArrayInputStream; -import java.io.DataInputStream; -import java.io.IOException; +import java.io.*; + import java.net.InetAddress; -import java.net.UnknownHostException; import java.nio.ByteBuffer; import com.google.common.base.Charsets; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import org.apache.cassandra.net.IVerbHandler; import org.apache.cassandra.net.Message; -import org.apache.cassandra.net.MessagingService; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.cassandra.net.*; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; +import static com.google.common.base.Charsets.UTF_8; + public class RowMutationVerbHandler implements IVerbHandler { @@ -67,11 +69,6 @@ public class RowMutationVerbHandler implements IVerbHandler hintedMutation.apply(); } } - - // Check if there were any forwarding headers in this message - byte[] forwardBytes = message.getHeader(RowMutation.FORWARD_HEADER); - if (forwardBytes != null) - forwardToLocalNodes(message, forwardBytes); Table.open(rm.getTable()).apply(rm, bytes, true); @@ -85,34 +82,5 @@ public class RowMutationVerbHandler implements IVerbHandler { logger_.error("Error in row mutation", e); } - } - - private void forwardToLocalNodes(Message message, byte[] forwardBytes) throws UnknownHostException - { - // remove fwds from message to avoid infinite loop - message.setHeader(RowMutation.FORWARD_HEADER, null); - - int bytesPerInetAddress = FBUtilities.getLocalAddress().getAddress().length; - assert forwardBytes.length >= bytesPerInetAddress; - assert forwardBytes.length % bytesPerInetAddress == 0; - - int offset = 0; - byte[] addressBytes = new byte[bytesPerInetAddress]; - - // Send a message to each of the addresses on our Forward List - while (offset < forwardBytes.length) - { - System.arraycopy(forwardBytes, offset, addressBytes, 0, bytesPerInetAddress); - InetAddress address = InetAddress.getByAddress(addressBytes); - - if (logger_.isDebugEnabled()) - logger_.debug("Forwarding message to " + address); - - // Send the original message to the address specified by the FORWARD_HINT - // Let the response go back to the coordinator - MessagingService.instance.sendOneWay(message, message.getFrom()); - - offset += bytesPerInetAddress; - } } } diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index cc652ae3ee..9b96607d55 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -27,23 +27,20 @@ import java.util.concurrent.*; import javax.management.MBeanServer; import javax.management.ObjectName; -import com.google.common.collect.HashMultimap; import com.google.common.collect.Multimap; +import static com.google.common.base.Charsets.UTF_8; import org.apache.commons.lang.ArrayUtils; import org.apache.commons.lang.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; + import org.apache.cassandra.concurrent.Stage; import org.apache.cassandra.concurrent.StageManager; import org.apache.cassandra.config.CFMetaData; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.*; -import org.apache.cassandra.db.filter.QueryFilter; -import org.apache.cassandra.dht.AbstractBounds; -import org.apache.cassandra.dht.Bounds; -import org.apache.cassandra.dht.IPartitioner; -import org.apache.cassandra.dht.Token; +import org.apache.cassandra.dht.*; import org.apache.cassandra.gms.Gossiper; import org.apache.cassandra.locator.AbstractReplicationStrategy; import org.apache.cassandra.locator.TokenMetadata; @@ -56,8 +53,7 @@ import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.LatencyTracker; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.WrappedRunnable; - -import static com.google.common.base.Charsets.UTF_8; +import org.apache.cassandra.db.filter.QueryFilter; public class StorageProxy implements StorageProxyMBean { @@ -94,14 +90,13 @@ public class StorageProxy implements StorageProxyMBean * @param mutations the mutations to be applied across the replicas * @param consistency_level the consistency level for the operation */ - public static void mutate(List mutations, ConsistencyLevel consistencyLevel) throws UnavailableException, TimeoutException + public static void mutate(List mutations, ConsistencyLevel consistency_level) throws UnavailableException, TimeoutException { long startTime = System.nanoTime(); - List responseHandlers = new ArrayList(); + ArrayList responseHandlers = new ArrayList(); RowMutation mostRecentRowMutation = null; StorageService ss = StorageService.instance; - String localDataCenter = getDataCenter(FBUtilities.getLocalAddress()); try { @@ -115,67 +110,58 @@ public class StorageProxy implements StorageProxyMBean Collection writeEndpoints = ss.getTokenMetadata().getWriteEndpoints(StorageService.getPartitioner().getToken(rm.key()), table, naturalEndpoints); Multimap hintedEndpoints = rs.getHintedEndpoints(writeEndpoints); - final IWriteResponseHandler responseHandler = rs.getWriteResponseHandler(writeEndpoints, hintedEndpoints, consistencyLevel); - - // exit early if we can't fulfuill the CL at this time + // send out the writes, as in mutate() above, but this time with a callback that tracks responses + final IWriteResponseHandler responseHandler = rs.getWriteResponseHandler(writeEndpoints, hintedEndpoints, consistency_level); responseHandler.assureSufficientLiveNodes(); - + responseHandlers.add(responseHandler); - - // Creates a Multimap that holds onto all the messages and addresses meant for a specific datacenter. - Multimap> dcMap = groupEndpointsByDataCenter(rm, hintedEndpoints, responseHandler); - - // Traverse all dataCenters where messages will be sent to. - for (Map.Entry>> entry : dcMap.asMap().entrySet()) + Message unhintedMessage = null; + for (Map.Entry> entry : hintedEndpoints.asMap().entrySet()) { - String dataCenter = entry.getKey(); - - // Grab a set of all the messages bound for this dataCenter and create an iterator over this set. - Collection> messagesForDataCenter = entry.getValue(); - Iterator> iter = messagesForDataCenter.iterator(); - assert iter.hasNext(); + InetAddress destination = entry.getKey(); + Collection targets = entry.getValue(); - // First endpoint in list is the destination for this group - Pair messageAndDestination = iter.next(); - - Message primaryMessage = messageAndDestination.left; - InetAddress target = messageAndDestination.right; - - // Add all the other destinations that are bound for the same dataCenter as a header in the primary message. - while (iter.hasNext()) + if (targets.size() == 1 && targets.iterator().next().equals(destination)) { - messageAndDestination = iter.next(); - assert messageAndDestination.left == primaryMessage; - - if (dataCenter.equals(localDataCenter)) + // unhinted writes + if (destination.equals(FBUtilities.getLocalAddress())) { - // direct write to local DC - assert primaryMessage.getHeader(RowMutation.FORWARD_HEADER) == null; - MessagingService.instance.sendOneWay(primaryMessage, target); + insertLocalMessage(rm, responseHandler); } else { - // group all nodes in this DC as forward headers on the primary message - ByteArrayOutputStream bos = new ByteArrayOutputStream(); - DataOutputStream dos = new DataOutputStream(bos); - - // append to older addresses - byte[] previousHints = primaryMessage.getHeader(RowMutation.FORWARD_HEADER); - if (previousHints != null) - dos.write(previousHints); - - dos.write(messageAndDestination.right.getAddress()); - primaryMessage.setHeader(RowMutation.FORWARD_HEADER, bos.toByteArray()); + // belongs on a different server. send it there. + if (unhintedMessage == null) + { + unhintedMessage = rm.makeRowMutationMessage(); + MessagingService.instance.addCallback(responseHandler, unhintedMessage.getMessageId()); + } + if (logger.isDebugEnabled()) + logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + unhintedMessage.getMessageId() + "@" + destination); + MessagingService.instance.sendOneWay(unhintedMessage, destination); } - } - - MessagingService.instance.sendOneWay(primaryMessage, target); + } + else + { + // hinted + Message hintedMessage = rm.makeRowMutationMessage(); + for (InetAddress target : targets) + { + if (!target.equals(destination)) + { + addHintHeader(hintedMessage, target); + if (logger.isDebugEnabled()) + logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + hintedMessage.getMessageId() + "@" + destination + " for " + target); + } + } + responseHandler.addHintCallback(hintedMessage, destination); + MessagingService.instance.sendOneWay(hintedMessage, destination); + } } } - // wait for writes. throws timeoutexception if necessary for (IWriteResponseHandler responseHandler : responseHandlers) - { + { responseHandler.get(); } } @@ -192,66 +178,6 @@ public class StorageProxy implements StorageProxyMBean } } - - private static Multimap> groupEndpointsByDataCenter(RowMutation rm, Multimap endpoints, final IWriteResponseHandler responseHandler) throws IOException - { - - Set>> endpointSet = endpoints.asMap().entrySet(); - Multimap> dcMap = HashMultimap.create(endpointSet.size(), 10); - Message unhintedMessage = null; - - for (Map.Entry> entry : endpointSet) - { - InetAddress destination = entry.getKey(); - Collection targets = entry.getValue(); - - String dataCenter = getDataCenter(destination); - - if (targets.size() == 1 && targets.iterator().next().equals(destination)) - { - // unhinted writes - if (destination.equals(FBUtilities.getLocalAddress())) - { - insertLocalMessage(rm, responseHandler); - } - else - { - // belongs on a different server. - if (unhintedMessage == null) - { - unhintedMessage = rm.makeRowMutationMessage(); - MessagingService.instance.addCallback(responseHandler, unhintedMessage.getMessageId()); - } - - if (logger.isDebugEnabled()) - logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + unhintedMessage.getMessageId() + "@" + destination); - - dcMap.put(dataCenter, new Pair(unhintedMessage, destination)); - } - } - else - { - // hinted - Message hintedMessage = rm.makeRowMutationMessage(); - - for (InetAddress target : targets) - { - if (!target.equals(destination)) - { - addHintHeader(hintedMessage, target); - if (logger.isDebugEnabled()) - logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + hintedMessage.getMessageId() + "@" + destination + " for " + target); - } - } - - responseHandler.addHintCallback(hintedMessage, destination); - dcMap.put(dataCenter, new Pair(hintedMessage, destination)); - } - } - - return dcMap; - } - private static void addHintHeader(Message message, InetAddress target) throws IOException { @@ -266,22 +192,6 @@ public class StorageProxy implements StorageProxyMBean message.setHeader(RowMutation.HINT, bos.toByteArray()); } - private static String getDataCenter(InetAddress addr) - { - String dataCenter = null; - try - { - dataCenter = DatabaseDescriptor.getEndpointSnitch().getDatacenter(addr); - } - catch (UnsupportedOperationException e) - { - // SimpleSnitch throws this - dataCenter = "default"; - } - - return dataCenter; - } - private static void insertLocalMessage(final RowMutation rm, final IWriteResponseHandler responseHandler) { if (logger.isDebugEnabled()) From d6b1d581e9ad28ea2106462d306304e821a335f9 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 2 Dec 2010 17:42:27 +0000 Subject: [PATCH 06/22] r/m out-of-date contrib/maven for CASSANDRA-1805 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041490 13f79535-47bb-0310-9956-ffa450edef68 --- contrib/maven/pom.xml | 551 ------------------------------------------ 1 file changed, 551 deletions(-) delete mode 100644 contrib/maven/pom.xml diff --git a/contrib/maven/pom.xml b/contrib/maven/pom.xml deleted file mode 100644 index 5d98eb76a8..0000000000 --- a/contrib/maven/pom.xml +++ /dev/null @@ -1,551 +0,0 @@ - - - - - org.apache - apache - 6 - - 4.0.0 - - org.apache.cassandra - cassandra - 0.6-SNAPSHOT - jar - Cassandra - 2009 - - - 2.0.9 - - - http://incubator.apache.org/cassandra - - - - cassandra-user - cassandra-user-subscribe@incubator.apache.org - cassandra-user-unsubscribe@incubator.apache.org - cassandra-user@incubator.apache.org - http://mail-archives.apache.org/mod_mbox/incubator-cassandra-user/ - - - cassandra-dev - cassandra-dev-subscribe@incubator.apache.org - cassandra-dev-unsubscribe@incubator.apache.org - cassandra-dev@incubator.apache.org - http://mail-archives.apache.org/mod_mbox/incubator-cassandra-dev/ - - - cassandra-commits - cassandra-commits-subscribe@incubator.apache.org - cassandra-commits-unsubscribe@incubator.apache.org - cassandra-commits@incubator.apache.org - http://mail-archives.apache.org/mod_mbox/incubator-cassandra-commits/ - - - - - - - - - - - JIRA - https://issues.apache.org/jira/browse/CASSANDRA - - - - - hudson - http://hudson.zones.apache.org/hudson - - - - mail - true - true - false - false -
cassandra-commits@incubator.apache.org
-
-
-
- - - - cassandra-website - scp://people.apache.org/x1/www/incubator.apache.org/cassandra/maven/${pom.version} - - - - - scm:svn:http://svn.apache.org/repos/asf/incubator/cassandra/trunk - scm:svn:https://svn.apache.org/repos/asf/incubator/cassandra/trunk - http://svn.apache.org/viewvc/incubator/cassandra/trunk/ - - - - - - commons-collections - commons-collections - 3.2.1 - - - commons-cli - commons-cli - 1.1 - - - commons-lang - commons-lang - 2.4 - - - jline - jline - 0.9.94 - - - log4j - log4j - 1.2.15 - - - - javax.jms - jms - - - com.sun.jmx - jmxri - - - com.sun.jdmk - jmxtools - - - - - org.slf4j - slf4j-api - 1.5.8 - - - org.slf4j - slf4j-log4j12 - 1.5.8 - - - org.antlr - antlr-runtime - 3.1.3 - - - com.google.collections - google-collections - 1.0-rc1 - - - - - - high-scale-lib - high-scale-lib - UNKNOWN - system - ${basedir}/lib/high-scale-lib.jar - - - flexjson - flexjson - 1.7 - system - ${basedir}/lib/flexjson-1.7.jar - - - libthrift - libthrift - UNKNOWN - system - ${basedir}/lib/libthrift-r820831.jar - - - jsonsimple - jsonsimple - UNKNOWN - system - ${basedir}/lib/json_simple-1.1.jar - - - com.reardencommerce - clhm - UNKNOWN - system - ${basedir}/lib/clhm-production.jar - - - - - junit - junit - 4.6 - test - - - - - - ${basedir}/src/java - ${basedir}/test/unit - build/classes - - - - ${basedir}/test/conf - - **/* - - - - ${basedir}/test/resources - - *.json - - - - - - - - - org.antlr - antlr3-maven-plugin - 3.1.3-1 - - - process-sources - - antlr - - - ${basedir}/src/java - - - - - - - - org.codehaus.mojo - build-helper-maven-plugin - 1.3 - - - add-source - generate-sources - - add-source - - - - ${basedir}/interface/gen-java - - - - - - - - - org.apache.maven.plugins - maven-compiler-plugin - - 1.6 - 1.6 - true - true - true - true - - - - - - org.apache.maven.plugins - maven-surefire-plugin - - - - storage-config - ${basedir}/test/conf - - - always - - **/TestRingCache.java - - - - - - - org.codehaus.mojo - cobertura-maven-plugin - 2.0 - - - - - - - org.apache.maven.plugins - maven-release-plugin - 2.0-beta-9 - - true - false - clean install - deploy - -Papache-release - - - - org.codehaus.mojo - ianal-maven-plugin - 1.0-alpha-1 - - - org.codehaus.mojo - rat-maven-plugin - 1.0-alpha-3 - - false - - - - org.apache.maven.plugins - maven-enforcer-plugin - - - validate - - enforce - - - - - [2.0.9,) - - - - - - - - org.codehaus.mojo - ianal-maven-plugin - - - - verify-legal-files - - - true - - - - - - - - - - - org.apache.maven.plugins - maven-jxr-plugin - - - org.apache.maven.plugins - maven-surefire-report-plugin - - - org.apache.maven.plugins - maven-pmd-plugin - - - org.codehaus.mojo - taglist-maven-plugin - - - org.apache.maven.plugins - maven-javadoc-plugin - - - http://java.sun.com/j2se/1.6.0/docs/api/ - http://logging.apache.org/log4j/docs/api/ - - - true - 900m - 1.6 - - - - org.codehaus.mojo - cobertura-maven-plugin - 2.2 - - - html - xml - - - - - - - - - - - thrift - - process-sources - - - - maven-antrun-plugin - - - process-sources - - - - - - - - run - - - - - - - - - - - apache-release - - - - - - org.apache.maven.plugins - maven-gpg-plugin - - ${gpg.passphrase} - - - - - sign - - - - - - - - true - org.apache.maven.plugins - maven-deploy-plugin - - true - - - - org.apache.maven.plugins - maven-source-plugin - - - attach-sources - - jar - - - - - - org.apache.maven.plugins - maven-javadoc-plugin - - ${project.build.sourceEncoding} - - - - attach-javadocs - - jar - - - - - - - - org.apache.maven.plugins - maven-assembly-plugin - - - - single - - package - - true - - - source-release - - - - - - - - org.apache.geronimo.genesis - apache-source-release-assembly-descriptor - 2.0 - - - - - - - - - - -
From 56770fc8070f15188e80e85978e235d75dc9f39f Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Dec 2010 14:23:44 +0000 Subject: [PATCH 07/22] fix cli case sensitivity for index queries patch by Pavel Yaskevich; reviewed by jbellis for CASSANDRA-1809' git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041833 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/cli/CliClient.java | 2 +- test/unit/org/apache/cassandra/cli/CliTest.java | 2 ++ 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/cli/CliClient.java b/src/java/org/apache/cassandra/cli/CliClient.java index 0aafe0c658..5f38de1ccb 100644 --- a/src/java/org/apache/cassandra/cli/CliClient.java +++ b/src/java/org/apache/cassandra/cli/CliClient.java @@ -468,7 +468,7 @@ public class CliClient extends CliUserHelp return; IndexClause clause = new IndexClause(); - String columnFamily = statement.getChild(0).getText(); + String columnFamily = CliCompiler.getColumnFamily(statement, keyspacesMap.get(keySpace).cf_defs); // ^(CONDITIONS ^(CONDITION $column $value) ...) Tree conditions = statement.getChild(1); diff --git a/test/unit/org/apache/cassandra/cli/CliTest.java b/test/unit/org/apache/cassandra/cli/CliTest.java index 50d4a01b7d..26f5082b23 100644 --- a/test/unit/org/apache/cassandra/cli/CliTest.java +++ b/test/unit/org/apache/cassandra/cli/CliTest.java @@ -41,6 +41,8 @@ public class CliTest extends CleanupHelper "get CF1[hello][world];", "set CF1[hello][world2] = 15;", "get CF1 where world2 = long(15);", + "get cF1 where world2 = long(15);", + "get Cf1 where world2 = long(15);", "set CF1['hello'][time_spent_uuid] = timeuuid(a8098c1a-f86e-11da-bd1a-00112444be1e);", "create column family CF2 with comparator=IntegerType;", "set CF2['key'][98349387493847748398334] = 'some text';", From 5856ee9e6e7732922bc2c331f4cdf7ee741bd661 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Dec 2010 14:25:24 +0000 Subject: [PATCH 08/22] cli support index type enum names patch by Pavel Yaskevich; reviewed by jbellis for CASSANDRA-1810 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041834 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../org/apache/cassandra/cli/CliClient.java | 26 ++++++++++++------- .../org/apache/cassandra/cli/CliTest.java | 2 +- 3 files changed, 19 insertions(+), 10 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 7342354997..3243f4f22c 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -25,6 +25,7 @@ dev * fix range queries against wrapped range (CASSANDRA-1781) * fix consistencylevel calculations for NetworkTopologyStrategy (CASSANDRA-1804) + * cli support index type enum names (CASSANDRA-1810) 0.7.0-rc1 diff --git a/src/java/org/apache/cassandra/cli/CliClient.java b/src/java/org/apache/cassandra/cli/CliClient.java index 5f38de1ccb..192ffa4ee9 100644 --- a/src/java/org/apache/cassandra/cli/CliClient.java +++ b/src/java/org/apache/cassandra/cli/CliClient.java @@ -1414,20 +1414,28 @@ public class CliClient extends CliUserHelp */ private IndexType getIndexTypeFromString(String indexTypeAsString) { - Integer indexTypeId; IndexType indexType; - try { - indexTypeId = new Integer(indexTypeAsString); + try + { + indexType = IndexType.findByValue(new Integer(indexTypeAsString)); } - catch (NumberFormatException e) { - throw new RuntimeException("Could not convert " + indexTypeAsString + " into Integer."); + catch (NumberFormatException e) + { + try + { + // if this is not an integer lets try to get IndexType by name + indexType = IndexType.valueOf(indexTypeAsString); + } + catch (IllegalArgumentException ie) + { + throw new RuntimeException("IndexType '" + indexTypeAsString + "' is unsupported."); + } } - indexType = IndexType.findByValue(indexTypeId); - - if (indexType == null) { - throw new RuntimeException(indexTypeAsString + " is unsupported."); + if (indexType == null) + { + throw new RuntimeException("IndexType '" + indexTypeAsString + "' is unsupported."); } return indexType; diff --git a/test/unit/org/apache/cassandra/cli/CliTest.java b/test/unit/org/apache/cassandra/cli/CliTest.java index 26f5082b23..6b74e78f35 100644 --- a/test/unit/org/apache/cassandra/cli/CliTest.java +++ b/test/unit/org/apache/cassandra/cli/CliTest.java @@ -36,7 +36,7 @@ public class CliTest extends CleanupHelper // please add new statements here so they could be auto-runned by this test. private String[] statements = { "use TestKeySpace;", - "create column family CF1 with comparator=UTF8Type and column_metadata=[{ column_name:world, validation_class:IntegerType, index_type:0, index_name:IdxName }, { column_name:world2, validation_class:LongType, index_type:0, index_name:LongIdxName}];", + "create column family CF1 with comparator=UTF8Type and column_metadata=[{ column_name:world, validation_class:IntegerType, index_type:0, index_name:IdxName }, { column_name:world2, validation_class:LongType, index_type:KEYS, index_name:LongIdxName}];", "set CF1[hello][world] = 123848374878933948398384;", "get CF1[hello][world];", "set CF1[hello][world2] = 15;", From 4ea000beb0561ad8dbd01cd44ac318ab5e6961a9 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Dec 2010 15:36:45 +0000 Subject: [PATCH 09/22] add debug messages for system_ thrift calls patch by jbellis git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041883 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/thrift/CassandraServer.java | 36 +++++++++---------- 1 file changed, 16 insertions(+), 20 deletions(-) diff --git a/src/java/org/apache/cassandra/thrift/CassandraServer.java b/src/java/org/apache/cassandra/thrift/CassandraServer.java index 028fb6e169..b7d0d40204 100644 --- a/src/java/org/apache/cassandra/thrift/CassandraServer.java +++ b/src/java/org/apache/cassandra/thrift/CassandraServer.java @@ -260,8 +260,7 @@ public class CassandraServer implements Cassandra.Iface public List get_slice(ByteBuffer key, ColumnParent column_parent, SlicePredicate predicate, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("get_slice"); + logger.debug("get_slice"); state().hasColumnFamilyAccess(column_parent.column_family, Permission.READ); return multigetSliceInternal(state().getKeyspace(), Collections.singletonList(key), column_parent, predicate, consistency_level).get(key); @@ -270,8 +269,7 @@ public class CassandraServer implements Cassandra.Iface public Map> multiget_slice(List keys, ColumnParent column_parent, SlicePredicate predicate, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("multiget_slice"); + logger.debug("multiget_slice"); state().hasColumnFamilyAccess(column_parent.column_family, Permission.READ); @@ -309,8 +307,7 @@ public class CassandraServer implements Cassandra.Iface public ColumnOrSuperColumn get(ByteBuffer key, ColumnPath column_path, ConsistencyLevel consistency_level) throws InvalidRequestException, NotFoundException, UnavailableException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("get"); + logger.debug("get"); state().hasColumnFamilyAccess(column_path.column_family, Permission.READ); String keyspace = state().getKeyspace(); @@ -338,8 +335,7 @@ public class CassandraServer implements Cassandra.Iface public int get_count(ByteBuffer key, ColumnParent column_parent, SlicePredicate predicate, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("get_count"); + logger.debug("get_count"); state().hasColumnFamilyAccess(column_parent.column_family, Permission.READ); @@ -349,8 +345,7 @@ public class CassandraServer implements Cassandra.Iface public Map multiget_count(List keys, ColumnParent column_parent, SlicePredicate predicate, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("multiget_count"); + logger.debug("multiget_count"); state().hasColumnFamilyAccess(column_parent.column_family, Permission.READ); String keyspace = state().getKeyspace(); @@ -367,8 +362,7 @@ public class CassandraServer implements Cassandra.Iface public void insert(ByteBuffer key, ColumnParent column_parent, Column column, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("insert"); + logger.debug("insert"); state().hasColumnFamilyAccess(column_parent.column_family, Permission.WRITE); @@ -391,8 +385,7 @@ public class CassandraServer implements Cassandra.Iface public void batch_mutate(Map>> mutation_map, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("batch_mutate"); + logger.debug("batch_mutate"); List cfamsSeen = new ArrayList(); @@ -428,8 +421,7 @@ public class CassandraServer implements Cassandra.Iface public void remove(ByteBuffer key, ColumnPath column_path, long timestamp, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("remove"); + logger.debug("remove"); state().hasColumnFamilyAccess(column_path.column_family, Permission.WRITE); @@ -482,8 +474,7 @@ public class CassandraServer implements Cassandra.Iface public List get_range_slices(ColumnParent column_parent, SlicePredicate predicate, KeyRange range, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TException, TimedOutException { - if (logger.isDebugEnabled()) - logger.debug("range_slice"); + logger.debug("range_slice"); String keyspace = state().getKeyspace(); state().hasColumnFamilyAccess(column_parent.column_family, Permission.READ); @@ -546,8 +537,7 @@ public class CassandraServer implements Cassandra.Iface public List get_indexed_slices(ColumnParent column_parent, IndexClause index_clause, SlicePredicate column_predicate, ConsistencyLevel consistency_level) throws InvalidRequestException, UnavailableException, TimedOutException, TException { - if (logger.isDebugEnabled()) - logger.debug("scan"); + logger.debug("scan"); state().hasColumnFamilyAccess(column_parent.column_family, Permission.READ); String keyspace = state().getKeyspace(); @@ -703,6 +693,7 @@ public class CassandraServer implements Cassandra.Iface public String system_add_column_family(CfDef cf_def) throws InvalidRequestException, TException { + logger.debug("add_column_family"); state().hasColumnFamilyListAccess(Permission.WRITE); ThriftValidation.validateCfDef(cf_def); try @@ -726,6 +717,7 @@ public class CassandraServer implements Cassandra.Iface public String system_drop_column_family(String column_family) throws InvalidRequestException, TException { + logger.debug("drop_column_family"); state().hasColumnFamilyListAccess(Permission.WRITE); try @@ -749,6 +741,7 @@ public class CassandraServer implements Cassandra.Iface public String system_add_keyspace(KsDef ks_def) throws InvalidRequestException, TException { + logger.debug("add_keyspace"); state().hasKeyspaceListAccess(Permission.WRITE); // generate a meaningful error if the user setup keyspace and/or column definition incorrectly @@ -792,6 +785,7 @@ public class CassandraServer implements Cassandra.Iface public String system_drop_keyspace(String keyspace) throws InvalidRequestException, TException { + logger.debug("drop_keyspace"); state().hasKeyspaceListAccess(Permission.WRITE); try @@ -816,6 +810,7 @@ public class CassandraServer implements Cassandra.Iface /** update an existing keyspace, but do not allow column family modifications. */ public String system_update_keyspace(KsDef ks_def) throws InvalidRequestException, TException { + logger.debug("update_keyspace"); state().hasKeyspaceListAccess(Permission.WRITE); ThriftValidation.validateTable(ks_def.name); @@ -848,6 +843,7 @@ public class CassandraServer implements Cassandra.Iface public String system_update_column_family(CfDef cf_def) throws InvalidRequestException, TException { + logger.debug("update_column_family"); state().hasColumnFamilyListAccess(Permission.WRITE); if (cf_def.keyspace == null || cf_def.name == null) From bf22c526c60d3f4bb2b23216af78dc178b4ac02e Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Dec 2010 15:41:25 +0000 Subject: [PATCH 10/22] validateCfDef on add_keyspace patch by jbellis git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041884 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/thrift/CassandraServer.java | 1 + 1 file changed, 1 insertion(+) diff --git a/src/java/org/apache/cassandra/thrift/CassandraServer.java b/src/java/org/apache/cassandra/thrift/CassandraServer.java index b7d0d40204..f0f72ff186 100644 --- a/src/java/org/apache/cassandra/thrift/CassandraServer.java +++ b/src/java/org/apache/cassandra/thrift/CassandraServer.java @@ -758,6 +758,7 @@ public class CassandraServer implements Cassandra.Iface Collection cfDefs = new ArrayList(ks_def.cf_defs.size()); for (CfDef cfDef : ks_def.cf_defs) { + ThriftValidation.validateCfDef(cfDef); cfDefs.add(convertToCFMetaData(cfDef)); } From 59b566bec3e0910287b18c9f6187cbeab63afd9a Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Dec 2010 17:52:46 +0000 Subject: [PATCH 11/22] improved validation of column_metadata patch by jbellis; reviewed by gdusbabek for CASSANDRA-1813 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041932 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + .../apache/cassandra/config/ColumnDefinition.java | 10 ---------- .../org/apache/cassandra/db/ColumnFamilyType.java | 5 +++-- .../apache/cassandra/thrift/ThriftValidation.java | 14 +++++++++++++- 4 files changed, 17 insertions(+), 13 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 3243f4f22c..571da4ccc4 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -26,6 +26,7 @@ dev * fix consistencylevel calculations for NetworkTopologyStrategy (CASSANDRA-1804) * cli support index type enum names (CASSANDRA-1810) + * improved validation of column_metadata (CASSANDRA-1813) 0.7.0-rc1 diff --git a/src/java/org/apache/cassandra/config/ColumnDefinition.java b/src/java/org/apache/cassandra/config/ColumnDefinition.java index 54f5b09ae6..20f574252f 100644 --- a/src/java/org/apache/cassandra/config/ColumnDefinition.java +++ b/src/java/org/apache/cassandra/config/ColumnDefinition.java @@ -103,7 +103,6 @@ public class ColumnDefinition { public static ColumnDefinition fromColumnDef(ColumnDef thriftColumnDef) throws ConfigurationException { - validateIndexType(thriftColumnDef); return new ColumnDefinition(thriftColumnDef.name, thriftColumnDef.validation_class, thriftColumnDef.index_type, thriftColumnDef.index_name); } @@ -123,10 +122,7 @@ public class ColumnDefinition { Map cds = new TreeMap(); for (ColumnDef thriftColumnDef : thriftDefs) - { - validateIndexType(thriftColumnDef); cds.put(thriftColumnDef.name, fromColumnDef(thriftColumnDef)); - } return Collections.unmodifiableMap(cds); } @@ -146,12 +142,6 @@ public class ColumnDefinition { return Collections.unmodifiableMap(cds); } - public static void validateIndexType(org.apache.cassandra.thrift.ColumnDef thriftColumnDef) throws ConfigurationException - { - if ((thriftColumnDef.index_name != null) && (thriftColumnDef.index_type == null)) - throw new ConfigurationException("index_name cannot be set if index_type is not also set"); - } - public static void validateIndexType(org.apache.cassandra.avro.ColumnDef avroColumnDef) throws ConfigurationException { if ((avroColumnDef.index_name != null) && (avroColumnDef.index_type == null)) diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyType.java b/src/java/org/apache/cassandra/db/ColumnFamilyType.java index 9f66b41346..d250ecd9b9 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyType.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyType.java @@ -25,11 +25,12 @@ public enum ColumnFamilyType Standard, Super; - public final static ColumnFamilyType create(String name) + public static ColumnFamilyType create(String name) { try { - return name == null ? null : ColumnFamilyType.valueOf(name); + // TODO thrift optional parameter in CfDef is leaking down here which it shouldn't + return name == null ? ColumnFamilyType.Standard : ColumnFamilyType.valueOf(name); } catch (IllegalArgumentException e) { diff --git a/src/java/org/apache/cassandra/thrift/ThriftValidation.java b/src/java/org/apache/cassandra/thrift/ThriftValidation.java index 57db71be41..062ba6af79 100644 --- a/src/java/org/apache/cassandra/thrift/ThriftValidation.java +++ b/src/java/org/apache/cassandra/thrift/ThriftValidation.java @@ -383,14 +383,20 @@ public class ThriftValidation { try { + ColumnFamilyType cfType = ColumnFamilyType.create(cf_def.column_type); + if (cfType == null) + throw new InvalidRequestException("invalid column type " + cf_def.column_type); + DatabaseDescriptor.getComparator(cf_def.comparator_type); DatabaseDescriptor.getComparator(cf_def.subcomparator_type); DatabaseDescriptor.getComparator(cf_def.default_validation_class); + if (cfType != ColumnFamilyType.Super && cf_def.subcomparator_type != null) + throw new InvalidRequestException("subcomparator_type is invalid for standard columns"); if (cf_def.column_metadata == null) return; - AbstractType comparator = cf_def.subcomparator_type == null + AbstractType comparator = cfType == ColumnFamilyType.Standard ? DatabaseDescriptor.getComparator(cf_def.comparator_type) : DatabaseDescriptor.getComparator(cf_def.subcomparator_type); for (ColumnDef c : cf_def.column_metadata) @@ -406,6 +412,12 @@ public class ThriftValidation throw new InvalidRequestException(String.format("Column name %s is not valid for comparator %s", FBUtilities.bytesToHex(c.name), cf_def.comparator_type)); } + + if ((c.index_name != null) && (c.index_type == null)) + throw new ConfigurationException("index_name cannot be set without index_type"); + + if (cfType == ColumnFamilyType.Super && c.index_type != null) + throw new InvalidRequestException("Secondary indexes are not supported on supercolumns"); } } catch (ConfigurationException e) From 2c6f56ddcb1248b14a27c6a1462de25c29177acf Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Dec 2010 18:52:06 +0000 Subject: [PATCH 12/22] reads at ConsistencyLevel > 1 throwUnavailableException immediately if insufficient live nodes exist patch by jbellis and tjake for CASSANDRA-1803 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041951 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 + .../cassandra/db/HintedHandOffManager.java | 3 +- .../cassandra/net/MessagingService.java | 14 +- .../cassandra/service/ConsistencyChecker.java | 7 +- .../DatacenterQuorumResponseHandler.java | 22 ++- .../cassandra/service/IResponseResolver.java | 6 +- .../service/QuorumResponseHandler.java | 43 +++--- .../service/RangeSliceResponseResolver.java | 19 ++- .../service/ReadResponseResolver.java | 36 +++-- .../cassandra/service/StorageProxy.java | 30 ++-- .../cassandra/thrift/CassandraServer.java | 8 +- .../service/ConsistencyLevelTest.java | 145 ++++++++++++++++++ 12 files changed, 262 insertions(+), 73 deletions(-) create mode 100644 test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java diff --git a/CHANGES.txt b/CHANGES.txt index 571da4ccc4..21563bbfeb 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -27,6 +27,8 @@ dev (CASSANDRA-1804) * cli support index type enum names (CASSANDRA-1810) * improved validation of column_metadata (CASSANDRA-1813) + * reads at ConsistencyLevel > 1 throw UnavailableException + immediately if insufficient live nodes exist (CASSANDRA-1803) 0.7.0-rc1 diff --git a/src/java/org/apache/cassandra/db/HintedHandOffManager.java b/src/java/org/apache/cassandra/db/HintedHandOffManager.java index dac514549f..c83bb646c3 100644 --- a/src/java/org/apache/cassandra/db/HintedHandOffManager.java +++ b/src/java/org/apache/cassandra/db/HintedHandOffManager.java @@ -24,6 +24,7 @@ import java.io.IOException; import java.net.InetAddress; import java.net.UnknownHostException; import java.nio.ByteBuffer; +import java.util.Arrays; import java.util.Collection; import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeoutException; @@ -122,7 +123,7 @@ public class HintedHandOffManager rm.add(cf); Message message = rm.makeRowMutationMessage(); IWriteResponseHandler responseHandler = WriteResponseHandler.create(endpoint); - MessagingService.instance.sendRR(message, new InetAddress[] { endpoint }, responseHandler); + MessagingService.instance.sendRR(message, Arrays.asList(endpoint), responseHandler); try { responseHandler.get(); diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index 5da901c3bb..883eb0a3f7 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -29,9 +29,7 @@ import java.nio.ByteBuffer; import java.nio.channels.AsynchronousCloseException; import java.nio.channels.ServerSocketChannel; import java.security.MessageDigest; -import java.util.EnumMap; -import java.util.HashMap; -import java.util.Map; +import java.util.*; import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -226,7 +224,7 @@ public class MessagingService implements MessagingServiceMBean * @return an reference to an IAsyncResult which can be queried for the * response */ - public String sendRR(Message message, InetAddress[] to, IAsyncCallback cb) + public String sendRR(Message message, Collection to, IAsyncCallback cb) { String messageId = message.getMessageId(); addCallback(cb, messageId); @@ -273,18 +271,16 @@ public class MessagingService implements MessagingServiceMBean * suggest that a timeout occured to the invoker of the send(). * @return an reference to message id used to match with the result */ - public String sendRR(Message[] messages, InetAddress[] to, IAsyncCallback cb) + public String sendRR(Message[] messages, List to, IAsyncCallback cb) { - if ( messages.length != to.length ) - { + if (messages.length != to.size()) throw new IllegalArgumentException("Number of messages and the number of endpoints need to be same."); - } String groupId = GuidGenerator.guid(); addCallback(cb, groupId); for ( int i = 0; i < messages.length; ++i ) { messages[i].setMessageId(groupId); - sendOneWay(messages[i], to[i]); + sendOneWay(messages[i], to.get(i)); } return groupId; } diff --git a/src/java/org/apache/cassandra/service/ConsistencyChecker.java b/src/java/org/apache/cassandra/service/ConsistencyChecker.java index a30b67ff93..f1bc2f2af7 100644 --- a/src/java/org/apache/cassandra/service/ConsistencyChecker.java +++ b/src/java/org/apache/cassandra/service/ConsistencyChecker.java @@ -156,7 +156,6 @@ class ConsistencyChecker implements Runnable static class DataRepairHandler implements IAsyncCallback { - private final Collection responses_ = new LinkedBlockingQueue(); private final ReadResponseResolver readResponseResolver_; private final int majority_; @@ -167,7 +166,6 @@ class ConsistencyChecker implements Runnable // wrap localRow in a response Message so it doesn't need to be special-cased in the resolver ReadResponse readResponse = new ReadResponse(localRow); Message fakeMessage = new Message(FBUtilities.getLocalAddress(), StorageService.Verb.INTERNAL_RESPONSE, ArrayUtils.EMPTY_BYTE_ARRAY); - responses_.add(fakeMessage); readResponseResolver_.injectPreProcessed(fakeMessage, readResponse); } @@ -176,15 +174,14 @@ class ConsistencyChecker implements Runnable { if (logger_.isDebugEnabled()) logger_.debug("Received response in DataRepairHandler : " + message.toString()); - responses_.add(message); readResponseResolver_.preprocess(message); - if (responses_.size() == majority_) + if (readResponseResolver_.getMessageCount() == majority_) { Runnable runnable = new WrappedRunnable() { public void runMayThrow() throws IOException, DigestMismatchException { - readResponseResolver_.resolve(responses_); + readResponseResolver_.resolve(); } }; // give remaining replicas until timeout to reply and get added to responses_ diff --git a/src/java/org/apache/cassandra/service/DatacenterQuorumResponseHandler.java b/src/java/org/apache/cassandra/service/DatacenterQuorumResponseHandler.java index 6267df94e9..0180d51dda 100644 --- a/src/java/org/apache/cassandra/service/DatacenterQuorumResponseHandler.java +++ b/src/java/org/apache/cassandra/service/DatacenterQuorumResponseHandler.java @@ -21,6 +21,9 @@ package org.apache.cassandra.service; */ +import java.net.InetAddress; +import java.util.Collection; +import java.util.List; import java.util.concurrent.atomic.AtomicInteger; import org.apache.cassandra.config.DatabaseDescriptor; @@ -29,6 +32,7 @@ import org.apache.cassandra.locator.IEndpointSnitch; import org.apache.cassandra.locator.NetworkTopologyStrategy; import org.apache.cassandra.net.Message; import org.apache.cassandra.thrift.ConsistencyLevel; +import org.apache.cassandra.thrift.UnavailableException; import org.apache.cassandra.utils.FBUtilities; /** @@ -49,14 +53,14 @@ public class DatacenterQuorumResponseHandler extends QuorumResponseHandler @Override public void response(Message message) { - responses.add(message); // we'll go ahead and resolve a reply from anyone, even if it's not from this dc + resolver.preprocess(message); int n; n = localdc.equals(snitch.getDatacenter(message.getFrom())) ? localResponses.decrementAndGet() : localResponses.get(); - if (n == 0 && responseResolver.isDataPresent(responses)) + if (n == 0 && resolver.isDataPresent()) { condition.signal(); } @@ -68,4 +72,18 @@ public class DatacenterQuorumResponseHandler extends QuorumResponseHandler NetworkTopologyStrategy stategy = (NetworkTopologyStrategy) Table.open(table).getReplicationStrategy(); return (stategy.getReplicationFactor(localdc) / 2) + 1; } + + @Override + public void assureSufficientLiveNodes(Collection endpoints) throws UnavailableException + { + int localEndpoints = 0; + for (InetAddress endpoint : endpoints) + { + if (localdc.equals(snitch.getDatacenter(endpoint))) + localEndpoints++; + } + + if(localEndpoints < blockfor) + throw new UnavailableException(); + } } diff --git a/src/java/org/apache/cassandra/service/IResponseResolver.java b/src/java/org/apache/cassandra/service/IResponseResolver.java index 0b5b54f64e..ea61705455 100644 --- a/src/java/org/apache/cassandra/service/IResponseResolver.java +++ b/src/java/org/apache/cassandra/service/IResponseResolver.java @@ -34,8 +34,10 @@ public interface IResponseResolver { * repairs . Hence you need to derive a response resolver based on your * needs from this interface. */ - public T resolve(Collection responses) throws DigestMismatchException, IOException; - public boolean isDataPresent(Collection responses); + public T resolve() throws DigestMismatchException, IOException; + public boolean isDataPresent(); public void preprocess(Message message); + public Iterable getMessages(); + public int getMessageCount(); } diff --git a/src/java/org/apache/cassandra/service/QuorumResponseHandler.java b/src/java/org/apache/cassandra/service/QuorumResponseHandler.java index a703e05718..9d7c7fdf02 100644 --- a/src/java/org/apache/cassandra/service/QuorumResponseHandler.java +++ b/src/java/org/apache/cassandra/service/QuorumResponseHandler.java @@ -18,11 +18,11 @@ package org.apache.cassandra.service; +import java.io.IOException; +import java.net.InetAddress; import java.util.Collection; -import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; -import java.io.IOException; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.Table; @@ -30,8 +30,8 @@ import org.apache.cassandra.net.IAsyncCallback; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.thrift.ConsistencyLevel; +import org.apache.cassandra.thrift.UnavailableException; import org.apache.cassandra.utils.SimpleCondition; - import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -39,19 +39,20 @@ public class QuorumResponseHandler implements IAsyncCallback { protected static final Logger logger = LoggerFactory.getLogger( QuorumResponseHandler.class ); protected final SimpleCondition condition = new SimpleCondition(); - protected final Collection responses = new LinkedBlockingQueue();; - protected IResponseResolver responseResolver; + protected final IResponseResolver resolver; private final long startTime; - protected int blockfor; + protected final int blockfor; /** * Constructor when response count has to be calculated and blocked for. */ - public QuorumResponseHandler(IResponseResolver responseResolver, ConsistencyLevel consistencyLevel, String table) + public QuorumResponseHandler(IResponseResolver resolver, ConsistencyLevel consistencyLevel, String table) { this.blockfor = determineBlockFor(consistencyLevel, table); - this.responseResolver = responseResolver; + this.resolver = resolver; this.startTime = System.currentTimeMillis(); + + logger.debug("QuorumResponseHandler blocking for {} responses", blockfor); } public T get() throws TimeoutException, DigestMismatchException, IOException @@ -72,35 +73,31 @@ public class QuorumResponseHandler implements IAsyncCallback if (!success) { StringBuilder sb = new StringBuilder(""); - for (Message message : responses) + for (Message message : resolver.getMessages()) { sb.append(message.getFrom()); } - throw new TimeoutException("Operation timed out - received only " + responses.size() + " responses from " + sb.toString() + " ."); + throw new TimeoutException("Operation timed out - received only " + resolver.getMessageCount() + " responses from " + sb.toString() + " ."); } } finally { - for (Message response : responses) + for (Message response : resolver.getMessages()) { MessagingService.removeRegisteredCallback(response.getMessageId()); } } - return responseResolver.resolve(responses); + return resolver.resolve(); } public void response(Message message) { - responses.add(message); - responseResolver.preprocess(message); - if (responses.size() < blockfor) { + resolver.preprocess(message); + if (resolver.getMessageCount() < blockfor) return; - } - if (responseResolver.isDataPresent(responses)) - { + if (resolver.isDataPresent()) condition.signal(); - } } public int determineBlockFor(ConsistencyLevel consistencyLevel, String table) @@ -115,7 +112,13 @@ public class QuorumResponseHandler implements IAsyncCallback case ALL: return Table.open(table).getReplicationStrategy().getReplicationFactor(); default: - throw new UnsupportedOperationException("invalid consistency level: " + table.toString()); + throw new UnsupportedOperationException("invalid consistency level: " + consistencyLevel); } } + + public void assureSufficientLiveNodes(Collection endpoints) throws UnavailableException + { + if (endpoints.size() < blockfor) + throw new UnavailableException(); + } } diff --git a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java index bec0c1296d..8b98019a74 100644 --- a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java +++ b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java @@ -21,6 +21,7 @@ package org.apache.cassandra.service; import java.io.IOException; import java.net.InetAddress; import java.util.*; +import java.util.concurrent.LinkedBlockingQueue; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -45,6 +46,7 @@ public class RangeSliceResponseResolver implements IResponseResolver> private static final Logger logger_ = LoggerFactory.getLogger(RangeSliceResponseResolver.class); private final String table; private final List sources; + protected final Collection responses = new LinkedBlockingQueue();; public RangeSliceResponseResolver(String table, List sources) { @@ -53,7 +55,7 @@ public class RangeSliceResponseResolver implements IResponseResolver> this.table = table; } - public List resolve(Collection responses) throws DigestMismatchException, IOException + public List resolve() throws DigestMismatchException, IOException { CollatingIterator collator = new CollatingIterator(new Comparator>() { @@ -110,11 +112,12 @@ public class RangeSliceResponseResolver implements IResponseResolver> public void preprocess(Message message) { + responses.add(message); } - public boolean isDataPresent(Collection responses) + public boolean isDataPresent() { - return responses.size() >= sources.size(); + return !responses.isEmpty(); } private static class RowIterator extends AbstractIterator> @@ -134,4 +137,14 @@ public class RangeSliceResponseResolver implements IResponseResolver> return iter.hasNext() ? new Pair(iter.next(), source) : endOfData(); } } + + public Iterable getMessages() + { + return responses; + } + + public int getMessageCount() + { + return responses.size(); + } } diff --git a/src/java/org/apache/cassandra/service/ReadResponseResolver.java b/src/java/org/apache/cassandra/service/ReadResponseResolver.java index 4a4d658b64..440f163776 100644 --- a/src/java/org/apache/cassandra/service/ReadResponseResolver.java +++ b/src/java/org/apache/cassandra/service/ReadResponseResolver.java @@ -58,14 +58,14 @@ public class ReadResponseResolver implements IResponseResolver * repair request should be scheduled. * */ - public Row resolve(Collection responses) throws DigestMismatchException, IOException + public Row resolve() throws DigestMismatchException, IOException { if (logger_.isDebugEnabled()) - logger_.debug("resolving " + responses.size() + " responses"); + logger_.debug("resolving " + results.size() + " responses"); long startTime = System.currentTimeMillis(); - List versions = new ArrayList(responses.size()); - List endpoints = new ArrayList(responses.size()); + List versions = new ArrayList(); + List endpoints = new ArrayList(); DecoratedKey key = null; ByteBuffer digest = FBUtilities.EMPTY_BYTE_BUFFER; boolean isDigestQuery = false; @@ -76,11 +76,10 @@ public class ReadResponseResolver implements IResponseResolver * query exists then we need to compare the digest with * the digest of the data that is received. */ - for (Message message : responses) - { - ReadResponse result = results.get(message); - if (result == null) - continue; // arrived after quorum already achieved + for (Map.Entry entry : results.entrySet()) + { + ReadResponse result = entry.getValue(); + Message message = entry.getKey(); if (result.isDigestQuery()) { digest = result.digest(); @@ -187,6 +186,8 @@ public class ReadResponseResolver implements IResponseResolver try { ReadResponse result = ReadResponse.serializer().deserialize(new DataInputStream(bufIn)); + if (logger_.isDebugEnabled()) + logger_.debug("Preprocessed {} response", result.isDigestQuery() ? "digest" : "data"); results.put(message, result); } catch (IOException e) @@ -201,16 +202,23 @@ public class ReadResponseResolver implements IResponseResolver results.put(message, result); } - public boolean isDataPresent(Collection responses) + public boolean isDataPresent() { - for (Message message : responses) + for (ReadResponse result : results.values()) { - ReadResponse result = results.get(message); - if (result == null) - continue; // arrived concurrently if (!result.isDigestQuery()) return true; } return false; } + + public Iterable getMessages() + { + return results.keySet(); + } + + public int getMessageCount() + { + return results.size(); + } } diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 9b96607d55..32a02dce0f 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -314,7 +314,7 @@ public class StorageProxy implements StorageProxyMBean private static List strongRead(List commands, ConsistencyLevel consistency_level) throws IOException, UnavailableException, TimeoutException { List> quorumResponseHandlers = new ArrayList>(); - List commandEndpoints = new ArrayList(); + List> commandEndpoints = new ArrayList>(); List rows = new ArrayList(); // send out read requests @@ -327,25 +327,25 @@ public class StorageProxy implements StorageProxyMBean Message messageDigestOnly = readMessageDigestOnly.makeReadMessage(); InetAddress dataPoint = StorageService.instance.findSuitableEndpoint(command.table, command.key); - List endpointList = StorageService.instance.getLiveNaturalEndpoints(command.table, command.key); + List endpoints = StorageService.instance.getLiveNaturalEndpoints(command.table, command.key); - InetAddress[] endpoints = new InetAddress[endpointList.size()]; - Message messages[] = new Message[endpointList.size()]; + AbstractReplicationStrategy rs = Table.open(command.table).getReplicationStrategy(); + QuorumResponseHandler handler = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), consistency_level); + handler.assureSufficientLiveNodes(endpoints); + + Message messages[] = new Message[endpoints.size()]; // data-request message is sent to dataPoint, the node that will actually get // the data for us. The other replicas are only sent a digest query. int n = 0; - for (InetAddress endpoint : endpointList) + for (InetAddress endpoint : endpoints) { Message m = endpoint.equals(dataPoint) ? message : messageDigestOnly; - endpoints[n] = endpoint; messages[n++] = m; if (logger.isDebugEnabled()) logger.debug("strongread reading " + (m == message ? "data" : "digest") + " for " + command + " from " + m.getMessageId() + "@" + endpoint); } - AbstractReplicationStrategy rs = Table.open(command.table).getReplicationStrategy(); - QuorumResponseHandler quorumResponseHandler = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), consistency_level); - MessagingService.instance.sendRR(messages, endpoints, quorumResponseHandler); - quorumResponseHandlers.add(quorumResponseHandler); + MessagingService.instance.sendRR(messages, endpoints, handler); + quorumResponseHandlers.add(handler); commandEndpoints.add(endpoints); } @@ -369,14 +369,14 @@ public class StorageProxy implements StorageProxyMBean catch (DigestMismatchException ex) { AbstractReplicationStrategy rs = Table.open(command.table).getReplicationStrategy(); - QuorumResponseHandler qrhRepair = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), ConsistencyLevel.QUORUM); + QuorumResponseHandler handler = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), ConsistencyLevel.QUORUM); if (logger.isDebugEnabled()) logger.debug("Digest mismatch:", ex); Message messageRepair = command.makeReadMessage(); - MessagingService.instance.sendRR(messageRepair, commandEndpoints.get(i), qrhRepair); + MessagingService.instance.sendRR(messageRepair, commandEndpoints.get(i), handler); if (repairResponseHandlers == null) repairResponseHandlers = new ArrayList>(); - repairResponseHandlers.add(qrhRepair); + repairResponseHandlers.add(handler); } } @@ -498,7 +498,7 @@ public class StorageProxy implements StorageProxyMBean final Message msg = new Message(FBUtilities.getLocalAddress(), StorageService.Verb.SCHEMA_CHECK, ArrayUtils.EMPTY_BYTE_ARRAY); final CountDownLatch latch = new CountDownLatch(liveHosts.size()); // an empty message acts as a request to the SchemaCheckVerbHandler. - MessagingService.instance.sendRR(msg, liveHosts.toArray(new InetAddress[]{}), new IAsyncCallback() + MessagingService.instance.sendRR(msg, liveHosts, new IAsyncCallback() { public void response(Message msg) { @@ -775,7 +775,7 @@ public class StorageProxy implements StorageProxyMBean logger.debug("Starting to send truncate messages to hosts {}", allEndpoints); Truncation truncation = new Truncation(keyspace, cfname); Message message = truncation.makeTruncationMessage(); - MessagingService.instance.sendRR(message, allEndpoints.toArray(new InetAddress[]{}), responseHandler); + MessagingService.instance.sendRR(message, allEndpoints, responseHandler); // Wait for all logger.debug("Sent all truncate messages, now waiting for {} responses", blockFor); diff --git a/src/java/org/apache/cassandra/thrift/CassandraServer.java b/src/java/org/apache/cassandra/thrift/CassandraServer.java index f0f72ff186..e8d659c465 100644 --- a/src/java/org/apache/cassandra/thrift/CassandraServer.java +++ b/src/java/org/apache/cassandra/thrift/CassandraServer.java @@ -138,6 +138,7 @@ public class CassandraServer implements Cassandra.Iface } catch (TimeoutException e) { + logger.debug("... timed out"); throw new TimedOutException(); } catch (IOException e) @@ -442,11 +443,12 @@ public class CassandraServer implements Cassandra.Iface try { - StorageProxy.mutate(mutations, consistency_level); + StorageProxy.mutate(mutations, consistency_level); } catch (TimeoutException e) { - throw new TimedOutException(); + logger.debug("... timed out"); + throw new TimedOutException(); } } finally @@ -512,6 +514,7 @@ public class CassandraServer implements Cassandra.Iface } catch (TimeoutException e) { + logger.debug("... timed out"); throw new TimedOutException(); } catch (IOException e) @@ -556,6 +559,7 @@ public class CassandraServer implements Cassandra.Iface } catch (TimeoutException e) { + logger.debug("... timed out"); throw new TimedOutException(); } return thriftifyKeySlices(rows, column_parent, column_predicate); diff --git a/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java b/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java new file mode 100644 index 0000000000..de86cf4f62 --- /dev/null +++ b/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java @@ -0,0 +1,145 @@ +package org.apache.cassandra.service; + +import java.net.InetAddress; +import java.util.ArrayList; +import java.util.List; + +import com.google.common.collect.HashMultimap; +import org.junit.Test; + +import org.apache.cassandra.CleanupHelper; +import org.apache.cassandra.Util; +import org.apache.cassandra.config.ConfigurationException; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.db.Row; +import org.apache.cassandra.dht.IPartitioner; +import org.apache.cassandra.dht.RandomPartitioner; +import org.apache.cassandra.dht.Token; +import org.apache.cassandra.locator.AbstractReplicationStrategy; +import org.apache.cassandra.locator.SimpleSnitch; +import org.apache.cassandra.locator.TokenMetadata; +import org.apache.cassandra.thrift.ConsistencyLevel; +import org.apache.cassandra.thrift.UnavailableException; + +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +public class ConsistencyLevelTest extends CleanupHelper +{ + @Test + public void testReadWriteConsistencyChecks() throws Exception + { + StorageService ss = StorageService.instance; + final int RING_SIZE = 3; + + TokenMetadata tmd = ss.getTokenMetadata(); + tmd.clearUnsafe(); + IPartitioner partitioner = new RandomPartitioner(); + + ss.setPartitionerUnsafe(partitioner); + + ArrayList endpointTokens = new ArrayList(); + ArrayList keyTokens = new ArrayList(); + List hosts = new ArrayList(); + + Util.createInitialRing(ss, partitioner, endpointTokens, keyTokens, hosts, RING_SIZE); + + HashMultimap hintedNodes = HashMultimap.create(); + + + AbstractReplicationStrategy strategy; + + for (String table : DatabaseDescriptor.getNonSystemTables()) + { + strategy = getStrategy(table, tmd); + StorageService.calculatePendingRanges(strategy, table); + int replicationFactor = strategy.getReplicationFactor(); + if (replicationFactor < 2) + continue; + + for (ConsistencyLevel c : ConsistencyLevel.values()) + { + + if (c == ConsistencyLevel.EACH_QUORUM || c == ConsistencyLevel.LOCAL_QUORUM) + continue; + + for (int i = 0; i < replicationFactor; i++) + { + hintedNodes.clear(); + + for (int j = 0; j < i; j++) + { + hintedNodes.put(hosts.get(j), hosts.get(j)); + } + + IWriteResponseHandler writeHandler = strategy.getWriteResponseHandler(hosts, hintedNodes, c); + + QuorumResponseHandler readHandler = strategy.getQuorumResponseHandler(new ReadResponseResolver(table), c); + + boolean isWriteUnavailable = false; + boolean isReadUnavailable = false; + try + { + writeHandler.assureSufficientLiveNodes(); + } + catch (UnavailableException e) + { + isWriteUnavailable = true; + } + + try + { + readHandler.assureSufficientLiveNodes(hintedNodes.asMap().keySet()); + } + catch (UnavailableException e) + { + isReadUnavailable = true; + } + + //these should always match (in this kind of test) + assertTrue(isWriteUnavailable == isReadUnavailable); + + switch (c) + { + case ALL: + if (isWriteUnavailable) + assertTrue(hintedNodes.size() < replicationFactor); + else + assertTrue(hintedNodes.size() >= replicationFactor); + + break; + case ONE: + case ANY: + if (isWriteUnavailable) + assertTrue(hintedNodes.size() == 0); + else + assertTrue(hintedNodes.size() > 0); + break; + case QUORUM: + if (isWriteUnavailable) + assertTrue(hintedNodes.size() < (replicationFactor / 2 + 1)); + else + assertTrue(hintedNodes.size() >= (replicationFactor / 2 + 1)); + break; + default: + fail("Unhandled CL: " + c); + + } + } + } + return; + } + + fail("Test requires at least one table with RF > 1"); + } + + private AbstractReplicationStrategy getStrategy(String table, TokenMetadata tmd) throws ConfigurationException + { + return AbstractReplicationStrategy.createReplicationStrategy(table, + "org.apache.cassandra.locator.SimpleStrategy", + tmd, + new SimpleSnitch(), + null); + } + +} From 69f95c15b753a45a2d9bf3942a395262d208853a Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Dec 2010 18:59:02 +0000 Subject: [PATCH 13/22] add commented-out JVM_OPTS lines for GC logging patch by mdennis; reviewed by jbellis for CASSANDRA-1807 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041953 13f79535-47bb-0310-9956-ffa450edef68 --- conf/cassandra-env.sh | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/conf/cassandra-env.sh b/conf/cassandra-env.sh index 9a946ad68d..2870fdb34d 100644 --- a/conf/cassandra-env.sh +++ b/conf/cassandra-env.sh @@ -78,7 +78,7 @@ if [ "`uname`" = "Linux" ] ; then JVM_OPTS="$JVM_OPTS -Xss128k" fi -# GC tuning options. +# GC tuning options JVM_OPTS="$JVM_OPTS -XX:+UseParNewGC" JVM_OPTS="$JVM_OPTS -XX:+UseConcMarkSweepGC" JVM_OPTS="$JVM_OPTS -XX:+CMSParallelRemarkEnabled" @@ -87,6 +87,14 @@ JVM_OPTS="$JVM_OPTS -XX:MaxTenuringThreshold=1" JVM_OPTS="$JVM_OPTS -XX:CMSInitiatingOccupancyFraction=75" JVM_OPTS="$JVM_OPTS -XX:+UseCMSInitiatingOccupancyOnly" +# GC logging options -- uncomment to enable +# JVM_OPTS="$JVM_OPTS -XX:+PrintGCDetails" +# JVM_OPTS="$JVM_OPTS -XX:+PrintGCTimeStamps" +# JVM_OPTS="$JVM_OPTS -XX:+PrintClassHistogram" +# JVM_OPTS="$JVM_OPTS -XX:+PrintTenuringDistribution" +# JVM_OPTS="$JVM_OPTS -XX:+PrintGCApplicationStoppedTime" +# JVM_OPTS="$JVM_OPTS -Xloggc:/var/log/cassandra/gc.log" + # Prefer binding to IPv4 network intefaces (when net.ipv6.bindv6only=1). See # http://bugs.sun.com/bugdatabase/view_bug.do?bug_id=6342561 (short version: # comment out this entry to enable IPv6 support). From 0b5d59ee5169e8f4e40c7b3c8c11f6f9e41ee4a4 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 3 Dec 2010 19:24:15 +0000 Subject: [PATCH 14/22] copy bytebuffers forlocal writes toavoid retainingthe entire Thrift frame patch by tjake; reviewed by jbellis for CASSANDRA-1801 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041961 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 ++ .../org/apache/cassandra/config/CFMetaData.java | 5 +++-- .../cassandra/config/ColumnDefinition.java | 5 +++-- src/java/org/apache/cassandra/db/Column.java | 6 ++++++ .../org/apache/cassandra/db/DeletedColumn.java | 7 +++++++ .../org/apache/cassandra/db/ExpiringColumn.java | 7 +++++++ src/java/org/apache/cassandra/db/IColumn.java | 3 +++ .../org/apache/cassandra/db/RowMutation.java | 16 ++++++++++++++++ .../org/apache/cassandra/db/SuperColumn.java | 16 ++++++++++++++++ .../apache/cassandra/service/StorageProxy.java | 2 +- .../apache/cassandra/utils/ByteBufferUtil.java | 14 ++++++++++++-- 11 files changed, 76 insertions(+), 7 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 21563bbfeb..185bf43e08 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -29,6 +29,8 @@ dev * improved validation of column_metadata (CASSANDRA-1813) * reads at ConsistencyLevel > 1 throw UnavailableException immediately if insufficient live nodes exist (CASSANDRA-1803) + * copy bytebuffers for local writes to avoid retaining the entire + Thrift frame (CASSANDRA-1801) 0.7.0-rc1 diff --git a/src/java/org/apache/cassandra/config/CFMetaData.java b/src/java/org/apache/cassandra/config/CFMetaData.java index 8f3f496858..59ec9f5a4e 100644 --- a/src/java/org/apache/cassandra/config/CFMetaData.java +++ b/src/java/org/apache/cassandra/config/CFMetaData.java @@ -41,6 +41,7 @@ import org.apache.cassandra.db.marshal.TimeUUIDType; import org.apache.cassandra.db.marshal.UTF8Type; import org.apache.cassandra.db.migration.Migration; import org.apache.cassandra.io.SerDeUtils; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.Pair; import org.apache.commons.lang.builder.EqualsBuilder; @@ -751,7 +752,7 @@ public final class CFMetaData org.apache.cassandra.avro.ColumnDef tcd = new org.apache.cassandra.avro.ColumnDef(); tcd.index_name = cd.getIndexName(); tcd.index_type = org.apache.cassandra.avro.IndexType.valueOf(cd.getIndexType().name()); - tcd.name = cd.name; + tcd.name = ByteBufferUtil.clone(cd.name); tcd.validation_class = cd.validator.getClass().getName(); column_meta.add(tcd); } @@ -786,7 +787,7 @@ public final class CFMetaData for (org.apache.cassandra.thrift.ColumnDef cdef : def.getColumn_metadata()) { org.apache.cassandra.avro.ColumnDef tdef = new org.apache.cassandra.avro.ColumnDef(); - tdef.name = cdef.BufferForName(); + tdef.name = ByteBufferUtil.clone(cdef.BufferForName()); tdef.validation_class = cdef.getValidation_class(); tdef.index_name = cdef.getIndex_name(); tdef.index_type = cdef.getIndex_type() == null ? null : org.apache.cassandra.avro.IndexType.valueOf(cdef.getIndex_type().name()); diff --git a/src/java/org/apache/cassandra/config/ColumnDefinition.java b/src/java/org/apache/cassandra/config/ColumnDefinition.java index 20f574252f..abb0c7ff75 100644 --- a/src/java/org/apache/cassandra/config/ColumnDefinition.java +++ b/src/java/org/apache/cassandra/config/ColumnDefinition.java @@ -31,6 +31,7 @@ import org.apache.avro.util.Utf8; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.thrift.ColumnDef; import org.apache.cassandra.thrift.IndexType; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; public class ColumnDefinition { @@ -103,7 +104,7 @@ public class ColumnDefinition { public static ColumnDefinition fromColumnDef(ColumnDef thriftColumnDef) throws ConfigurationException { - return new ColumnDefinition(thriftColumnDef.name, thriftColumnDef.validation_class, thriftColumnDef.index_type, thriftColumnDef.index_name); + return new ColumnDefinition(ByteBufferUtil.clone(thriftColumnDef.name), thriftColumnDef.validation_class, thriftColumnDef.index_type, thriftColumnDef.index_name); } public static ColumnDefinition fromColumnDef(org.apache.cassandra.avro.ColumnDef avroColumnDef) throws ConfigurationException @@ -122,7 +123,7 @@ public class ColumnDefinition { Map cds = new TreeMap(); for (ColumnDef thriftColumnDef : thriftDefs) - cds.put(thriftColumnDef.name, fromColumnDef(thriftColumnDef)); + cds.put(ByteBufferUtil.clone(thriftColumnDef.name), fromColumnDef(thriftColumnDef)); return Collections.unmodifiableMap(cds); } diff --git a/src/java/org/apache/cassandra/db/Column.java b/src/java/org/apache/cassandra/db/Column.java index b61ee42f49..38abb152f5 100644 --- a/src/java/org/apache/cassandra/db/Column.java +++ b/src/java/org/apache/cassandra/db/Column.java @@ -210,6 +210,12 @@ public class Column implements IColumn return result; } + @Override + public IColumn deepCopy() + { + return new Column(ByteBufferUtil.clone(name), ByteBufferUtil.clone(value), timestamp); + } + public String getString(AbstractType comparator) { StringBuilder sb = new StringBuilder(); diff --git a/src/java/org/apache/cassandra/db/DeletedColumn.java b/src/java/org/apache/cassandra/db/DeletedColumn.java index 5e7c7ec3f6..7af5af5190 100644 --- a/src/java/org/apache/cassandra/db/DeletedColumn.java +++ b/src/java/org/apache/cassandra/db/DeletedColumn.java @@ -20,6 +20,7 @@ package org.apache.cassandra.db; import java.nio.ByteBuffer; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -55,4 +56,10 @@ public class DeletedColumn extends Column { return value.getInt(value.position()+value.arrayOffset() ); } + + @Override + public IColumn deepCopy() + { + return new DeletedColumn(ByteBufferUtil.clone(name), ByteBufferUtil.clone(value), timestamp); + } } diff --git a/src/java/org/apache/cassandra/db/ExpiringColumn.java b/src/java/org/apache/cassandra/db/ExpiringColumn.java index 5496a45e3e..e587618ef6 100644 --- a/src/java/org/apache/cassandra/db/ExpiringColumn.java +++ b/src/java/org/apache/cassandra/db/ExpiringColumn.java @@ -24,6 +24,7 @@ import java.security.MessageDigest; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.io.util.DataOutputBuffer; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.log4j.Logger; /** @@ -104,6 +105,12 @@ public class ExpiringColumn extends Column return localExpirationTime; } + @Override + public IColumn deepCopy() + { + return new ExpiringColumn(ByteBufferUtil.clone(name), ByteBufferUtil.clone(value), timestamp, timeToLive, localExpirationTime); + } + @Override public String getString(AbstractType comparator) { diff --git a/src/java/org/apache/cassandra/db/IColumn.java b/src/java/org/apache/cassandra/db/IColumn.java index d7f9d3d642..61bb6783c3 100644 --- a/src/java/org/apache/cassandra/db/IColumn.java +++ b/src/java/org/apache/cassandra/db/IColumn.java @@ -46,6 +46,9 @@ public interface IColumn public int getLocalDeletionTime(); // for tombstone GC, so int is sufficient granularity public String getString(AbstractType comparator); + /** clones the column, making copies of any underlying byte buffers */ + IColumn deepCopy(); + /** * For a simple column, live == !isMarkedForDelete. * For a supercolumn, live means it has at least one subcolumn whose timestamp is greater than the diff --git a/src/java/org/apache/cassandra/db/RowMutation.java b/src/java/org/apache/cassandra/db/RowMutation.java index 5217d83ac0..4f59f3cc44 100644 --- a/src/java/org/apache/cassandra/db/RowMutation.java +++ b/src/java/org/apache/cassandra/db/RowMutation.java @@ -40,6 +40,7 @@ import org.apache.cassandra.service.StorageService; import org.apache.cassandra.thrift.ColumnOrSuperColumn; import org.apache.cassandra.thrift.Deletion; import org.apache.cassandra.thrift.Mutation; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; import org.apache.commons.lang.StringUtils; @@ -330,6 +331,21 @@ public class RowMutation rm.delete(new QueryPath(cfName, del.super_column), del.timestamp); } } + + public RowMutation deepCopy() + { + RowMutation rm = new RowMutation(table_, ByteBufferUtil.clone(key_)); + + for (Map.Entry e : modifications_.entrySet()) + { + ColumnFamily cf = e.getValue().cloneMeShallow(); + for (Map.Entry ce : e.getValue().getColumnsMap().entrySet()) + cf.addColumn(ce.getValue().deepCopy()); + rm.modifications_.put(e.getKey(), cf); + } + + return rm; + } } class RowMutationSerializer implements ICompactSerializer diff --git a/src/java/org/apache/cassandra/db/SuperColumn.java b/src/java/org/apache/cassandra/db/SuperColumn.java index c74bf371fd..0d5def4485 100644 --- a/src/java/org/apache/cassandra/db/SuperColumn.java +++ b/src/java/org/apache/cassandra/db/SuperColumn.java @@ -24,6 +24,7 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.security.MessageDigest; import java.util.Collection; +import java.util.Map; import java.util.concurrent.ConcurrentSkipListMap; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; @@ -31,6 +32,7 @@ import java.util.concurrent.atomic.AtomicLong; import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.io.ICompactSerializer2; import org.apache.cassandra.io.util.DataOutputBuffer; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -286,6 +288,20 @@ public class SuperColumn implements IColumn, IColumnContainer this.localDeletionTime.set(localDeleteTime); this.markedForDeleteAt.set(timestamp); } + + public IColumn deepCopy() + { + SuperColumn sc = new SuperColumn(ByteBufferUtil.clone(name_), this.getComparator()); + sc.localDeletionTime = localDeletionTime; + sc.markedForDeleteAt = markedForDeleteAt; + + for(Map.Entry c : columns_.entrySet()) + { + sc.addColumn(c.getValue().deepCopy()); + } + + return sc; + } public IColumn reconcile(IColumn c) { diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 32a02dce0f..e5020ad5e2 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -200,7 +200,7 @@ public class StorageProxy implements StorageProxyMBean { public void runMayThrow() throws IOException { - rm.apply(); + rm.deepCopy().apply(); responseHandler.response(null); } }; diff --git a/src/java/org/apache/cassandra/utils/ByteBufferUtil.java b/src/java/org/apache/cassandra/utils/ByteBufferUtil.java index a96a33001f..80fdd229ef 100644 --- a/src/java/org/apache/cassandra/utils/ByteBufferUtil.java +++ b/src/java/org/apache/cassandra/utils/ByteBufferUtil.java @@ -60,8 +60,8 @@ import java.nio.charset.Charset; * } * */ -public class ByteBufferUtil { - +public class ByteBufferUtil +{ public static int compareUnsigned(ByteBuffer o1, ByteBuffer o2) { return FBUtilities.compareUnsigned(o1.array(), o2.array(), o1.arrayOffset()+o1.position(), o2.arrayOffset()+o2.position(), o1.limit()+o1.arrayOffset(), o2.limit()+o2.arrayOffset()); @@ -98,4 +98,14 @@ public class ByteBufferUtil { throw new RuntimeException(e); } } + + public static ByteBuffer clone(ByteBuffer o) + { + ByteBuffer clone = ByteBuffer.allocate(o.remaining()); + o.mark(); + clone.put(o); + o.reset(); + clone.flip(); + return clone; + } } From 91eea35343c1a6347f8ae98739ce2f6deff83ed9 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Fri, 3 Dec 2010 20:57:03 +0000 Subject: [PATCH 15/22] Reduce FatClient timeout to RING_DELAY / 2. Patch by brandonwilliams, reviewed by jbellis for CASSANDRA-1730. git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1041998 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + src/java/org/apache/cassandra/gms/Gossiper.java | 4 ++-- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 185bf43e08..72a535fb81 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -31,6 +31,7 @@ dev immediately if insufficient live nodes exist (CASSANDRA-1803) * copy bytebuffers for local writes to avoid retaining the entire Thrift frame (CASSANDRA-1801) + * reduce fat client timeout (CASSANDRA-1730) 0.7.0-rc1 diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index c79b9e30fb..dcc5095a3c 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -145,8 +145,8 @@ public class Gossiper implements IFailureDetectionEventListener { // 3 days aVeryLongTime_ = 259200 * 1000; - // 1 hour - FatClientTimeout_ = 60 * 60 * 1000; + // half of RING_DELAY, to ensure justRemovedEndpoints has enough leeway to prevent re-gossip + FatClientTimeout_ = (long)(StorageService.RING_DELAY / 2); /* register with the Failure Detector for receiving Failure detector events */ FailureDetector.instance.registerFailureDetectionEventListener(this); } From 93539ffd2cfe2dd14297bc8b5c13620a6cfcad0d Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Sat, 4 Dec 2010 05:00:46 +0000 Subject: [PATCH 16/22] fix botched merge of CASSANDRA-1316 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1042100 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 ++ src/java/org/apache/cassandra/service/StorageProxy.java | 2 +- 2 files changed, 3 insertions(+), 1 deletion(-) diff --git a/CHANGES.txt b/CHANGES.txt index 72a535fb81..6d362ecfc1 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -31,7 +31,9 @@ dev immediately if insufficient live nodes exist (CASSANDRA-1803) * copy bytebuffers for local writes to avoid retaining the entire Thrift frame (CASSANDRA-1801) + * fix NPE adding index to column w/o prior metadata (CASSANDRA-1764) * reduce fat client timeout (CASSANDRA-1730) + * fix botched merge of CASSANDRA-1316 0.7.0-rc1 diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index e5020ad5e2..08be299ca4 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -369,7 +369,7 @@ public class StorageProxy implements StorageProxyMBean catch (DigestMismatchException ex) { AbstractReplicationStrategy rs = Table.open(command.table).getReplicationStrategy(); - QuorumResponseHandler handler = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), ConsistencyLevel.QUORUM); + QuorumResponseHandler handler = rs.getQuorumResponseHandler(new ReadResponseResolver(command.table), consistency_level); if (logger.isDebugEnabled()) logger.debug("Digest mismatch:", ex); Message messageRepair = command.makeReadMessage(); From f3f93d80a0f3d6ab2ce6d7a990cec57549b38bb6 Mon Sep 17 00:00:00 2001 From: Eric Evans Date: Mon, 6 Dec 2010 17:22:31 +0000 Subject: [PATCH 17/22] update versioning for 0.7 rc2 release git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1042730 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 6d362ecfc1..f75075e126 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,4 @@ -dev +0.7.0-rc2 * fix live-column-count of slice ranges including tombstoned supercolumn with live subcolumn (CASSANDRA-1591) * rename o.a.c.internal.AntientropyStage -> AntiEntropyStage, diff --git a/build.xml b/build.xml index d4c9ab2f29..1e92f7f8a8 100644 --- a/build.xml +++ b/build.xml @@ -47,7 +47,7 @@ - + diff --git a/debian/changelog b/debian/changelog index 0017e47f0d..765ccda88f 100644 --- a/debian/changelog +++ b/debian/changelog @@ -1,3 +1,9 @@ +cassandra (0.7.0~rc2) unstable; urgency=low + + * Release candidate release. + + -- Eric Evans Mon, 06 Dec 2010 11:19:40 -0600 + cassandra (0.7.0~rc1) unstable; urgency=low * Release candidate release. From b52b98548d7a7e4e7b99d1c8e720dcd2fd9878fb Mon Sep 17 00:00:00 2001 From: Eric Evans Date: Mon, 6 Dec 2010 17:22:36 +0000 Subject: [PATCH 18/22] prepend missing license blurb git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1042731 13f79535-47bb-0310-9956-ffa450edef68 --- .../service/ConsistencyLevelTest.java | 21 +++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java b/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java index de86cf4f62..2ebb5bf395 100644 --- a/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java +++ b/test/unit/org/apache/cassandra/service/ConsistencyLevelTest.java @@ -1,4 +1,25 @@ package org.apache.cassandra.service; +/* + * + * 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. + * + */ + import java.net.InetAddress; import java.util.ArrayList; From 9cb32b628881f88df5dbfbf8b2001cca4523cd70 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Mon, 6 Dec 2010 23:24:17 +0000 Subject: [PATCH 19/22] Switch word_count CFs to AsciiType. Patch by brandonwillliams git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1042851 13f79535-47bb-0310-9956-ffa450edef68 --- contrib/word_count/src/WordCountSetup.java | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/contrib/word_count/src/WordCountSetup.java b/contrib/word_count/src/WordCountSetup.java index cfe1812b8d..5a6e00273b 100644 --- a/contrib/word_count/src/WordCountSetup.java +++ b/contrib/word_count/src/WordCountSetup.java @@ -99,8 +99,14 @@ public class WordCountSetup private static void setupKeyspace(Cassandra.Iface client) throws TException, InvalidRequestException { List cfDefList = new ArrayList(); - cfDefList.add(new CfDef(WordCount.KEYSPACE, WordCount.COLUMN_FAMILY)); - cfDefList.add(new CfDef(WordCount.KEYSPACE, WordCount.OUTPUT_COLUMN_FAMILY)); + CfDef input = new CfDef(WordCount.KEYSPACE, WordCount.COLUMN_FAMILY); + input.setComparator_type("AsciiType"); + input.setDefault_validation_class("AsciiType"); + cfDefList.add(input); + CfDef output = new CfDef(WordCount.KEYSPACE, WordCount.OUTPUT_COLUMN_FAMILY); + output.setComparator_type("AsciiType"); + output.setDefault_validation_class("AsciiType"); + cfDefList.add(output); client.system_add_keyspace(new KsDef(WordCount.KEYSPACE, "org.apache.cassandra.locator.SimpleStrategy", 1, cfDefList)); int magnitude = client.describe_ring(WordCount.KEYSPACE).size(); From 0aa85e9398065fd604012303fbc194055d0065f4 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Mon, 6 Dec 2010 23:34:05 +0000 Subject: [PATCH 20/22] word_count uses better ks/cf names now that we're not piggybacking off storage definitions in the config. Patch by brandonwilliams git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1042857 13f79535-47bb-0310-9956-ffa450edef68 --- contrib/word_count/src/WordCount.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/contrib/word_count/src/WordCount.java b/contrib/word_count/src/WordCount.java index 93aba94cb9..fe52d5f6dd 100644 --- a/contrib/word_count/src/WordCount.java +++ b/contrib/word_count/src/WordCount.java @@ -54,11 +54,11 @@ public class WordCount extends Configured implements Tool { private static final Logger logger = LoggerFactory.getLogger(WordCount.class); - static final String KEYSPACE = "Keyspace1"; - static final String COLUMN_FAMILY = "Standard1"; + static final String KEYSPACE = "wordcount"; + static final String COLUMN_FAMILY = "input_words"; static final String OUTPUT_REDUCER_VAR = "output_reducer"; - static final String OUTPUT_COLUMN_FAMILY = "Standard2"; + static final String OUTPUT_COLUMN_FAMILY = "output_words"; private static final String OUTPUT_PATH_PREFIX = "/tmp/word_count"; private static final String CONF_COLUMN_NAME = "columnname"; From 7eae748cc3ee7ca584c9d459f533b2dedc9e7375 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Mon, 6 Dec 2010 23:46:44 +0000 Subject: [PATCH 21/22] Update word_count README git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1042863 13f79535-47bb-0310-9956-ffa450edef68 --- contrib/word_count/README.txt | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/contrib/word_count/README.txt b/contrib/word_count/README.txt index 5fd57a0a69..338b7074c4 100644 --- a/contrib/word_count/README.txt +++ b/contrib/word_count/README.txt @@ -13,15 +13,15 @@ contrib/word_count$ bin/word_count The output of the word count can now be configured. In the bin/word_count file, you can specify the OUTPUT_REDUCER. The two options are 'filesystem' and 'cassandra'. The filesystem option outputs to the /tmp/word_count* -directories. The cassandra option outputs to the 'Standard2' column family. +directories. The cassandra option outputs to the 'output_words' column family +in the 'wordcount' keyspace. -In order to view the results in Cassandra, one can use python/pycassa and +In order to view the results in Cassandra, one can use bin/cassandra-cli and perform the following operations: -$ python ->>> import pycassa ->>> con = pycassa.connect('Keyspace1') ->>> cf = pycassa.ColumnFamily(con, 'Standard2') ->>> list(cf.get_range()) +$ bin/cassandra-cli +> connect localhost/9160 +> use wordcount; +> list output_words; Read the code in src/ for more details. From 3a63f3e2e00abd7d9c5adcdb185208a96c35a800 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Tue, 7 Dec 2010 21:06:53 +0000 Subject: [PATCH 22/22] nodetool can display compaction stats. Patch by Edward Capriolo, reviewed by brandonwilliams for CASSANDRA-1763 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1043201 13f79535-47bb-0310-9956-ffa450edef68 --- .../org/apache/cassandra/tools/NodeCmd.java | 19 +++++++++++++++++-- .../org/apache/cassandra/tools/NodeProbe.java | 10 ++++++++++ 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/src/java/org/apache/cassandra/tools/NodeCmd.java b/src/java/org/apache/cassandra/tools/NodeCmd.java index fbb65ad865..c2f969b80c 100644 --- a/src/java/org/apache/cassandra/tools/NodeCmd.java +++ b/src/java/org/apache/cassandra/tools/NodeCmd.java @@ -37,6 +37,7 @@ import org.apache.cassandra.cache.JMXInstrumentedCacheMBean; import org.apache.cassandra.concurrent.IExecutorMBean; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ColumnFamilyStoreMBean; +import org.apache.cassandra.db.CompactionManagerMBean; import org.apache.cassandra.dht.Token; import org.apache.cassandra.net.MessagingServiceMBean; @@ -75,7 +76,7 @@ public class NodeCmd { "clearsnapshot, tpstats, flush, drain, repair, decommission, move, loadbalance, removetoken [status|force]|[token], " + "setcachecapacity [keyspace] [cfname] [keycachecapacity] [rowcachecapacity], " + "getcompactionthreshold [keyspace] [cfname], setcompactionthreshold [cfname] [minthreshold] [maxthreshold], " + - "netstats [host], cfhistograms "); + "netstats [host], cfhistograms , compactionstats"); String usage = String.format("java %s --host %n", NodeCmd.class.getName()); hf.printHelp(usage, "", options, header); } @@ -247,7 +248,17 @@ public class NodeCmd { completed += n; outs.printf("%-25s%10s%10s%15s%n", "Responses", "n/a", pending, completed); } - + + public void printCompactionStats(PrintStream outs) + { + CompactionManagerMBean cm = probe.getCompactionManagerProxy(); + outs.println("compaction type: " + (cm.getCompactionType() == null ? "n/a" : cm.getCompactionType())); + outs.println("column family: " + (cm.getColumnFamilyInProgress() == null ? "n/a" : cm.getColumnFamilyInProgress())); + outs.println("bytes compacted: " + (cm.getBytesCompacted() == null ? "n/a" : cm.getBytesCompacted())); + outs.println("bytes total in progress: " + (cm.getBytesTotalInProgress() == null ? "n/a" : cm.getBytesTotalInProgress() )); + outs.println("pending tasks: " + cm.getPendingTasks()); + } + public void printColumnFamilyStats(PrintStream outs) { Map > cfstoreMap = new HashMap >(); @@ -493,6 +504,10 @@ public class NodeCmd { System.exit(3); } } + else if (cmdName.equals("compactionstats")) + { + nodeCmd.printCompactionStats(System.out); + } else if (cmdName.equals("cfstats")) { nodeCmd.printColumnFamilyStats(System.out); diff --git a/src/java/org/apache/cassandra/tools/NodeProbe.java b/src/java/org/apache/cassandra/tools/NodeProbe.java index 75ad89d5fb..50f9592b35 100644 --- a/src/java/org/apache/cassandra/tools/NodeProbe.java +++ b/src/java/org/apache/cassandra/tools/NodeProbe.java @@ -45,6 +45,8 @@ import org.apache.cassandra.cache.JMXInstrumentedCacheMBean; import org.apache.cassandra.concurrent.IExecutorMBean; import org.apache.cassandra.config.ConfigurationException; import org.apache.cassandra.db.ColumnFamilyStoreMBean; +import org.apache.cassandra.db.CompactionManager; +import org.apache.cassandra.db.CompactionManagerMBean; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.dht.Token; import org.apache.cassandra.locator.IEndpointSnitch; @@ -72,6 +74,7 @@ public class NodeProbe private JMXConnector jmxc; private MBeanServerConnection mbeanServerConn; + private CompactionManagerMBean compactionProxy; private StorageServiceMBean ssProxy; private MemoryMXBean memProxy; private RuntimeMXBean runtimeProxy; @@ -121,6 +124,8 @@ public class NodeProbe ssProxy = JMX.newMBeanProxy(mbeanServerConn, name, StorageServiceMBean.class); name = new ObjectName(StreamingService.MBEAN_OBJECT_NAME); streamProxy = JMX.newMBeanProxy(mbeanServerConn, name, StreamingServiceMBean.class); + name = new ObjectName(CompactionManager.MBEAN_OBJECT_NAME); + compactionProxy = JMX.newMBeanProxy(mbeanServerConn, name, CompactionManagerMBean.class); } catch (MalformedObjectNameException e) { throw new RuntimeException( @@ -224,6 +229,11 @@ public class NodeProbe } } + public CompactionManagerMBean getCompactionManagerProxy() + { + return compactionProxy; + } + public JMXInstrumentedCacheMBean getKeyCacheMBean(String tableName, String cfName) { String keyCachePath = "org.apache.cassandra.db:type=Caches,keyspace=" + tableName + ",cache=" + cfName + "KeyCache";