From c67a22ffe4fc2a86aa43b4d770889b8d3c84f988 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 20 Jan 2011 02:34:07 +0000 Subject: [PATCH 01/15] fix messages/endpoints mismatch when RR is disabled patch by jbellis; reviewed by tjake for CASSANDRA-2010 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061104 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/service/StorageProxy.java | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index c106ddd28a..6b78762dd9 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -354,18 +354,17 @@ public class StorageProxy implements StorageProxyMBean ReadCallback handler = getReadCallback(resolver, command.table, consistency_level); handler.assureSufficientLiveNodes(endpoints); - int targets; + // if we're not going to read repair, cut the endpoints list down to the ones required to satisfy ConsistencyLevel if (randomlyReadRepair(command)) { - targets = endpoints.size(); - if (targets > handler.blockfor) + if (endpoints.size() > handler.blockfor) repairs.add(command); } else { - targets = handler.blockfor; + endpoints = endpoints.subList(0, handler.blockfor); } - Message[] messages = new Message[targets]; + 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. From a6fa6013651c7702261425f30503871fe0fffdab Mon Sep 17 00:00:00 2001 From: Eric Evans Date: Thu, 20 Jan 2011 15:56:56 +0000 Subject: [PATCH 02/15] added missing license headers Patch eevans; reported by Stephen Connolly for CASSANDRA-2016 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061353 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/io/util/ColumnSortedMap.java | 23 ++++++++++++++++++- .../cassandra/service/RepairCallback.java | 21 +++++++++++++++++ .../utils/BloomFilterSerializer.java | 21 +++++++++++++++++ .../utils/LegacyBloomFilterSerializer.java | 21 +++++++++++++++++ 4 files changed, 85 insertions(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/io/util/ColumnSortedMap.java b/src/java/org/apache/cassandra/io/util/ColumnSortedMap.java index 98f7abe668..76aae7a386 100644 --- a/src/java/org/apache/cassandra/io/util/ColumnSortedMap.java +++ b/src/java/org/apache/cassandra/io/util/ColumnSortedMap.java @@ -1,4 +1,25 @@ package org.apache.cassandra.io.util; +/* + * + * 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.io.DataInput; import java.io.IOException; @@ -261,4 +282,4 @@ class ColumnIterator implements Iterator> { throw new UnsupportedOperationException(); } -} \ No newline at end of file +} diff --git a/src/java/org/apache/cassandra/service/RepairCallback.java b/src/java/org/apache/cassandra/service/RepairCallback.java index 7c485baba7..8ddd4849c3 100644 --- a/src/java/org/apache/cassandra/service/RepairCallback.java +++ b/src/java/org/apache/cassandra/service/RepairCallback.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.io.IOException; import java.net.InetAddress; diff --git a/src/java/org/apache/cassandra/utils/BloomFilterSerializer.java b/src/java/org/apache/cassandra/utils/BloomFilterSerializer.java index 5acdaf08a0..ad59e7c257 100644 --- a/src/java/org/apache/cassandra/utils/BloomFilterSerializer.java +++ b/src/java/org/apache/cassandra/utils/BloomFilterSerializer.java @@ -1,4 +1,25 @@ package org.apache.cassandra.utils; +/* + * + * 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.io.DataInputStream; import java.io.DataOutputStream; diff --git a/src/java/org/apache/cassandra/utils/LegacyBloomFilterSerializer.java b/src/java/org/apache/cassandra/utils/LegacyBloomFilterSerializer.java index a4add6916d..62b2df62e1 100644 --- a/src/java/org/apache/cassandra/utils/LegacyBloomFilterSerializer.java +++ b/src/java/org/apache/cassandra/utils/LegacyBloomFilterSerializer.java @@ -1,4 +1,25 @@ package org.apache.cassandra.utils; +/* + * + * 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.util.BitSet; import java.io.DataInputStream; From 0985590d4c9a1b57fc9a0a1b286d664a5ad2bb6e Mon Sep 17 00:00:00 2001 From: Eric Evans Date: Thu, 20 Jan 2011 16:11:02 +0000 Subject: [PATCH 03/15] chdir / on startup Patch by eevans for CASSANDRA-1718 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061364 13f79535-47bb-0310-9956-ffa450edef68 --- debian/init | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/debian/init b/debian/init index 1d2961fb51..c1c3707868 100644 --- a/debian/init +++ b/debian/init @@ -119,6 +119,9 @@ do_start() # 2 if daemon could not be started is_running && return 1 + cassandra_home=`getent passwd cassandra | awk -F ':' '{ print $6; }'` + cd / # jsvc doesn't chdir() for us + $JSVC \ -user cassandra \ -home $JAVA_HOME \ @@ -127,6 +130,8 @@ do_start() -outfile /var/log/$NAME/output.log \ -cp `classpath` \ -Dlog4j.configuration=log4j-server.properties \ + -XX:HeapDumpPath="$cassandra_home/java_`date +%s`.hprof" \ + -XX:ErrorFile="$cassandra_home/hs_err_`date +%s`.log" \ $JVM_OPTS \ org.apache.cassandra.thrift.CassandraDaemon From 6f6cd971a4a7c0595403751f8e1dea4e7e5effe4 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 20 Jan 2011 19:36:21 +0000 Subject: [PATCH 04/15] fix stress.java for FBU/BBU refactor git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061475 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/contrib/stress/tests/IndexedRangeSlicer.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/IndexedRangeSlicer.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/IndexedRangeSlicer.java index 333946b426..ec353d31cc 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/IndexedRangeSlicer.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/IndexedRangeSlicer.java @@ -19,6 +19,7 @@ package org.apache.cassandra.contrib.stress.tests; import org.apache.cassandra.contrib.stress.util.OperationThread; import org.apache.cassandra.thrift.*; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; import java.nio.ByteBuffer; @@ -102,11 +103,11 @@ public class IndexedRangeSlicer extends OperationThread private int getMaxKey(List keySlices) { byte[] firstKey = keySlices.get(0).getKey(); - int maxKey = FBUtilities.byteBufferToInt(ByteBuffer.wrap(firstKey)); + int maxKey = ByteBufferUtil.toInt(ByteBuffer.wrap(firstKey)); for (KeySlice k : keySlices) { - int currentKey = FBUtilities.byteBufferToInt(ByteBuffer.wrap(k.getKey())); + int currentKey = ByteBufferUtil.toInt(ByteBuffer.wrap(k.getKey())); if (currentKey > maxKey) { From 281cbbffec5df6ad9793dae5544995f9d58c53bd Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 20 Jan 2011 19:46:50 +0000 Subject: [PATCH 05/15] add stress.java README git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061477 13f79535-47bb-0310-9956-ffa450edef68 --- contrib/stress/README.txt | 59 +++++++++++++++++++ contrib/stress/bin/stress | 2 +- .../cassandra/contrib/stress/Session.java | 2 +- 3 files changed, 61 insertions(+), 2 deletions(-) create mode 100644 contrib/stress/README.txt diff --git a/contrib/stress/README.txt b/contrib/stress/README.txt new file mode 100644 index 0000000000..76f04a4c9f --- /dev/null +++ b/contrib/stress/README.txt @@ -0,0 +1,59 @@ +stress +====== + +Description +----------- +stress is a tool for benchmarking and load testing a Cassandra +cluster. It is significantly faster than the older py_stress tool. + +Setup +----- +Run `ant` from the Cassandra source directory, then Run `ant` from the +contrib/stress directory. + +Usage +----- +There are three different modes of operation: + + * inserting (loading test data) + * reading + * range slicing (only works with the OrderPreservingPartioner) + * indexed range slicing (works with RandomParitioner on indexed ColumnFamilies) + +Important options: + -o or --operation: + Sets the operation mode, one of 'insert', 'read', 'rangeslice', or 'indexedrangeslice' + -n or --num-keys: + the number of rows to insert/read/slice; defaults to one million + -d or --nodes: + the node(s) to perform the test against. For multiple nodes, supply a + comma-separated list without spaces, ex: cassandra1,cassandra2,cassandra3 + -y or --family-type: + Sets the ColumnFamily type. One of 'Standard' or 'Super'. If using super, + you probably want to set the -u option also. + -c or --columns: + the number of columns per row, defaults to 5 + -u or --supercolumns: + use the number of supercolumns specified NOTE: you must set the -y + option appropriately, or this option has no effect. + -g or --get-range-slice-count: + This is only used for the rangeslice operation and will *NOT* work with + the RandomPartioner. You must set the OrderPreservingPartioner in your + storage-conf.xml (note that you will need to wipe all existing data + when switching partioners.) This option sets the number of rows to + slice at a time and defaults to 1000. + -r or --random: + Only used for reads. By default, stress.py will perform reads on rows + with a guassian distribution, which will cause some repeats. Setting + this option makes the reads completely random instead. + -i or --progress-interval: + The interval, in seconds, at which progress will be output. + +Remember that you must perform inserts before performing reads or range slices. + +Examples +-------- + + * contrib/stress/bin/stress -d 192.168.1.101 # 1M inserts to given host + * contrib/stress/bin/stress -d 192.168.1.101 -o read # 1M reads + * contrib/stress/bin/stress -d 192.168.1.101,192.168.1.102 -n 10000000 # 10M inserts spread across two nodes diff --git a/contrib/stress/bin/stress b/contrib/stress/bin/stress index 1284b3c09a..d43e319348 100755 --- a/contrib/stress/bin/stress +++ b/contrib/stress/bin/stress @@ -23,7 +23,7 @@ if [ "x$CLASSPATH" = "x" ]; then exit 1 fi - # Circuit class files. + # Stress class files. if [ ! -d `dirname $0`/../build/classes ]; then echo "Unable to locate stress class files" >&2 exit 1 diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java index 56071f3c93..a29d868744 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java @@ -52,7 +52,7 @@ public class Session availableOptions.addOption("c", "columns", true, "Number of columns per key, default:5."); availableOptions.addOption("S", "column-size", true, "Size of column values in bytes, default:34."); availableOptions.addOption("C", "cardinality", true, "Number of unique values stored in columns, default:50."); - availableOptions.addOption("d", "nodes", true, "Host nodes (comma separated), default:locahost."); + availableOptions.addOption("d", "nodes", true, "Host nodes (--comma separated), default:locahost."); availableOptions.addOption("s", "stdev", true, "Standard Deviation Factor, default:0.1."); availableOptions.addOption("r", "random", false, "Use random key generator (STDEV will have no effect), default:false."); availableOptions.addOption("f", "file", true, "Write output to file"); From ef34db673735def1b4a52f1b3c623a5a4e674ddb Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 20 Jan 2011 20:07:07 +0000 Subject: [PATCH 06/15] clean up help output git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061484 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/contrib/stress/Session.java | 44 +++++++++---------- 1 file changed, 22 insertions(+), 22 deletions(-) diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java index a29d868744..853d5987fa 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/Session.java @@ -45,28 +45,28 @@ public class Session static { - availableOptions.addOption("h", "help", false, "show this help message and exit."); - availableOptions.addOption("n", "num-keys", true, "Number of keys, default:1000000."); - availableOptions.addOption("N", "skip-keys", true, "Fraction of keys to skip initially, default:0."); - availableOptions.addOption("t", "threads", true, "Number of threads to use, default:50."); - availableOptions.addOption("c", "columns", true, "Number of columns per key, default:5."); - availableOptions.addOption("S", "column-size", true, "Size of column values in bytes, default:34."); - availableOptions.addOption("C", "cardinality", true, "Number of unique values stored in columns, default:50."); - availableOptions.addOption("d", "nodes", true, "Host nodes (--comma separated), default:locahost."); - availableOptions.addOption("s", "stdev", true, "Standard Deviation Factor, default:0.1."); - availableOptions.addOption("r", "random", false, "Use random key generator (STDEV will have no effect), default:false."); - availableOptions.addOption("f", "file", true, "Write output to file"); - availableOptions.addOption("p", "port", true, "Thrift port, default:9160."); - availableOptions.addOption("m", "unframed", false, "Use unframed transport, default:false."); - availableOptions.addOption("o", "operation", true, "Operation to perform (INSERT, READ, RANGE_SLICE, INDEXED_RANGE_SLICE, MULTI_GET), default:INSERT."); - availableOptions.addOption("u", "supercolumns", true, "Number of super columns per key, default:1."); - availableOptions.addOption("y", "family-type", true, "Column Family Type (Super, Standard), default:Standard."); - availableOptions.addOption("k", "keep-going", false, "Ignore errors inserting or reading, default:false."); - availableOptions.addOption("i", "progress-interval", true, "Progress Report Interval (seconds), default:10."); - availableOptions.addOption("g", "keys-per-call", true, "Amount of keys to get_range_slices or multiget per call, default:1000."); - availableOptions.addOption("l", "replication-factor", true, "Replication Factor to use when creating needed column families, default:1."); - availableOptions.addOption("e", "consistency-level", true, "Consistency Level to use (ONE, QUORUM, LOCAL_QUORUM, EACH_QUORUM, ALL, ANY), default:ONE."); - availableOptions.addOption("x", "create-index", true, "Type of index to create on needed column families (KEYS)."); + availableOptions.addOption("h", "help", false, "Show this help message and exit"); + availableOptions.addOption("n", "num-keys", true, "Number of keys, default:1000000"); + availableOptions.addOption("N", "skip-keys", true, "Fraction of keys to skip initially, default:0"); + availableOptions.addOption("t", "threads", true, "Number of threads to use, default:50"); + availableOptions.addOption("c", "columns", true, "Number of columns per key, default:5"); + availableOptions.addOption("S", "column-size", true, "Size of column values in bytes, default:34"); + availableOptions.addOption("C", "cardinality", true, "Number of unique values stored in columns, default:50"); + availableOptions.addOption("d", "nodes", true, "Host nodes (comma separated), default:locahost"); + availableOptions.addOption("s", "stdev", true, "Standard Deviation Factor, default:0.1"); + availableOptions.addOption("r", "random", false, "Use random key generator (STDEV will have no effect), default:false"); + availableOptions.addOption("f", "file", true, "Write output to given file"); + availableOptions.addOption("p", "port", true, "Thrift port, default:9160"); + availableOptions.addOption("m", "unframed", false, "Use unframed transport, default:false"); + availableOptions.addOption("o", "operation", true, "Operation to perform (INSERT, READ, RANGE_SLICE, INDEXED_RANGE_SLICE, MULTI_GET), default:INSERT"); + availableOptions.addOption("u", "supercolumns", true, "Number of super columns per key, default:1"); + availableOptions.addOption("y", "family-type", true, "Column Family Type (Super, Standard), default:Standard"); + availableOptions.addOption("k", "keep-going", false, "Ignore errors inserting or reading, default:false"); + availableOptions.addOption("i", "progress-interval", true, "Progress Report Interval (seconds), default:10"); + availableOptions.addOption("g", "keys-per-call", true, "Number of keys to get_range_slices or multiget per call, default:1000"); + availableOptions.addOption("l", "replication-factor", true, "Replication Factor to use when creating needed column families, default:1"); + availableOptions.addOption("e", "consistency-level", true, "Consistency Level to use (ONE, QUORUM, LOCAL_QUORUM, EACH_QUORUM, ALL, ANY), default:ONE"); + availableOptions.addOption("x", "create-index", true, "Type of index to create on needed column families (KEYS)"); } private int numKeys = 1000 * 1000; From 2a3fee59ae9b44711d720de9503bace14d360cf7 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 20 Jan 2011 20:32:25 +0000 Subject: [PATCH 07/15] rename tests -> operations git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061497 13f79535-47bb-0310-9956-ffa450edef68 --- .../stress/src/org/apache/cassandra/contrib/stress/Stress.java | 2 +- .../stress/{tests => operations}/IndexedRangeSlicer.java | 2 +- .../contrib/stress/{tests => operations}/Inserter.java | 2 +- .../contrib/stress/{tests => operations}/MultiGetter.java | 2 +- .../contrib/stress/{tests => operations}/RangeSlicer.java | 2 +- .../cassandra/contrib/stress/{tests => operations}/Reader.java | 2 +- 6 files changed, 6 insertions(+), 6 deletions(-) rename contrib/stress/src/org/apache/cassandra/contrib/stress/{tests => operations}/IndexedRangeSlicer.java (98%) rename contrib/stress/src/org/apache/cassandra/contrib/stress/{tests => operations}/Inserter.java (98%) rename contrib/stress/src/org/apache/cassandra/contrib/stress/{tests => operations}/MultiGetter.java (98%) rename contrib/stress/src/org/apache/cassandra/contrib/stress/{tests => operations}/RangeSlicer.java (98%) rename contrib/stress/src/org/apache/cassandra/contrib/stress/{tests => operations}/Reader.java (98%) diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/Stress.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/Stress.java index 9ef1e4f077..72e0942322 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/Stress.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/Stress.java @@ -17,7 +17,7 @@ */ package org.apache.cassandra.contrib.stress; -import org.apache.cassandra.contrib.stress.tests.*; +import org.apache.cassandra.contrib.stress.operations.*; import org.apache.cassandra.contrib.stress.util.OperationThread; import org.apache.commons.cli.Option; diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/IndexedRangeSlicer.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/IndexedRangeSlicer.java similarity index 98% rename from contrib/stress/src/org/apache/cassandra/contrib/stress/tests/IndexedRangeSlicer.java rename to contrib/stress/src/org/apache/cassandra/contrib/stress/operations/IndexedRangeSlicer.java index ec353d31cc..6891921fe9 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/IndexedRangeSlicer.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/IndexedRangeSlicer.java @@ -15,7 +15,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.cassandra.contrib.stress.tests; +package org.apache.cassandra.contrib.stress.operations; import org.apache.cassandra.contrib.stress.util.OperationThread; import org.apache.cassandra.thrift.*; diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/Inserter.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Inserter.java similarity index 98% rename from contrib/stress/src/org/apache/cassandra/contrib/stress/tests/Inserter.java rename to contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Inserter.java index 6a9e48909f..7065b115c7 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/Inserter.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Inserter.java @@ -15,7 +15,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.cassandra.contrib.stress.tests; +package org.apache.cassandra.contrib.stress.operations; import org.apache.cassandra.contrib.stress.util.OperationThread; import org.apache.cassandra.db.ColumnFamilyType; diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/MultiGetter.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/MultiGetter.java similarity index 98% rename from contrib/stress/src/org/apache/cassandra/contrib/stress/tests/MultiGetter.java rename to contrib/stress/src/org/apache/cassandra/contrib/stress/operations/MultiGetter.java index dcead126ca..f0bbe76666 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/MultiGetter.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/MultiGetter.java @@ -15,7 +15,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.cassandra.contrib.stress.tests; +package org.apache.cassandra.contrib.stress.operations; import org.apache.cassandra.contrib.stress.util.OperationThread; import org.apache.cassandra.db.ColumnFamilyType; diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/RangeSlicer.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/RangeSlicer.java similarity index 98% rename from contrib/stress/src/org/apache/cassandra/contrib/stress/tests/RangeSlicer.java rename to contrib/stress/src/org/apache/cassandra/contrib/stress/operations/RangeSlicer.java index 54674acbcb..8cdf6d6f8b 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/RangeSlicer.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/RangeSlicer.java @@ -15,7 +15,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.cassandra.contrib.stress.tests; +package org.apache.cassandra.contrib.stress.operations; import org.apache.cassandra.contrib.stress.util.OperationThread; import org.apache.cassandra.db.ColumnFamilyType; diff --git a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/Reader.java b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Reader.java similarity index 98% rename from contrib/stress/src/org/apache/cassandra/contrib/stress/tests/Reader.java rename to contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Reader.java index 41e5d91dda..14552568aa 100644 --- a/contrib/stress/src/org/apache/cassandra/contrib/stress/tests/Reader.java +++ b/contrib/stress/src/org/apache/cassandra/contrib/stress/operations/Reader.java @@ -15,7 +15,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.cassandra.contrib.stress.tests; +package org.apache.cassandra.contrib.stress.operations; import org.apache.cassandra.contrib.stress.util.OperationThread; import org.apache.cassandra.db.ColumnFamilyType; From 6d96344145001efea019dacd96b555d91c494037 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 20 Jan 2011 21:31:26 +0000 Subject: [PATCH 08/15] s/currentNode/currentOwner/ git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061525 13f79535-47bb-0310-9956-ffa450edef68 --- .../org/apache/cassandra/service/StorageService.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 2b0fe93079..8f52f6dfdd 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -622,22 +622,22 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe logger_.info("Node " + endpoint + " state jump to normal"); // we don't want to update if this node is responsible for the token and it has a later startup time than endpoint. - InetAddress currentNode = tokenMetadata_.getEndpoint(token); - if (currentNode == null) + InetAddress currentOwner = tokenMetadata_.getEndpoint(token); + if (currentOwner == null) { logger_.debug("New node " + endpoint + " at token " + token); tokenMetadata_.updateNormalToken(token, endpoint); if (!isClientMode) SystemTable.updateToken(endpoint, token); } - else if (endpoint.equals(currentNode)) + else if (endpoint.equals(currentOwner)) { // nothing to do } - else if (Gossiper.instance.compareEndpointStartup(endpoint, currentNode) > 0) + else if (Gossiper.instance.compareEndpointStartup(endpoint, currentOwner) > 0) { logger_.info(String.format("Nodes %s and %s have the same token %s. %s is the new owner", - endpoint, currentNode, token, endpoint)); + endpoint, currentOwner, token, endpoint)); tokenMetadata_.updateNormalToken(token, endpoint); if (!isClientMode) SystemTable.updateToken(endpoint, token); @@ -645,7 +645,7 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe else { logger_.info(String.format("Nodes %s and %s have the same token %s. Ignoring %s", - endpoint, currentNode, token, endpoint)); + endpoint, currentOwner, token, endpoint)); } if (pieces.length > 2) From 96e530594a65b2b09abda9ad0716324a39222de6 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Thu, 20 Jan 2011 21:36:49 +0000 Subject: [PATCH 09/15] Update token metadata for NORMAL state when endpoint has not changed. Patch by brandonwilliams, reviewed by jbellis for CASSANDRA-1934 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061534 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/service/StorageService.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 8f52f6dfdd..19adf65498 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -632,7 +632,9 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe } else if (endpoint.equals(currentOwner)) { - // nothing to do + // set state back to normal, since the node may have tried to leave, but failed and is now back up + // no need to persist, token/ip did not change + tokenMetadata_.updateNormalToken(token, endpoint); } else if (Gossiper.instance.compareEndpointStartup(endpoint, currentOwner) > 0) { From 60b848b43f016091c89c77eb1713ea563136dbf2 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Thu, 20 Jan 2011 22:53:11 +0000 Subject: [PATCH 10/15] Add a configurable maximum amount of time to hint for a dead host. Patch by brandonwilliams, reviewed by jbellis for CASSANDRA-1459 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061557 13f79535-47bb-0310-9956-ffa450edef68 --- conf/cassandra.yaml | 4 ++++ .../org/apache/cassandra/config/Config.java | 1 + .../cassandra/config/DatabaseDescriptor.java | 5 +++++ .../org/apache/cassandra/gms/Gossiper.java | 19 +++++++++++++----- .../locator/AbstractReplicationStrategy.java | 7 +++++++ .../cassandra/service/StorageProxy.java | 20 +++++++++++++++++-- .../cassandra/service/StorageProxyMBean.java | 2 ++ 7 files changed, 51 insertions(+), 7 deletions(-) diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index 8241c78c0d..ca5495d5b9 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -31,6 +31,10 @@ auto_bootstrap: false # See http://wiki.apache.org/cassandra/HintedHandoff hinted_handoff_enabled: true +# this defines the maximum amount of time a dead host will have hints +# generated. After it has been dead this long, hints will be dropped. +# Maximum is approximately 50 days +max_hint_window_in_ms: 2147483647 # authentication backend, implementing IAuthenticator; used to identify users authenticator: org.apache.cassandra.auth.AllowAllAuthenticator diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index def0a5e04b..fd4bca0dbe 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -34,6 +34,7 @@ public class Config public Boolean auto_bootstrap = false; public Boolean hinted_handoff_enabled = true; + public Integer max_hint_window_in_ms = Integer.MAX_VALUE; public String[] seeds; public DiskAccessMode disk_access_mode = DiskAccessMode.auto; diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index f1f23f07ad..972cd5ea57 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -1079,6 +1079,11 @@ public class DatabaseDescriptor return conf.hinted_handoff_enabled; } + public static int getMaxHintWindow() + { + return conf.max_hint_window_in_ms; + } + public static AbstractType getValueValidator(String keyspace, String cf, ByteBuffer column) { return getCFMetaData(keyspace, cf).getValueValidator(column); diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index de863f831b..312cf85ae5 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -128,7 +128,7 @@ public class Gossiper implements IFailureDetectionEventListener private Set liveEndpoints_ = new ConcurrentSkipListSet(inetcomparator); /* unreachable member set */ - private Set unreachableEndpoints_ = new ConcurrentSkipListSet(inetcomparator); + private Map unreachableEndpoints_ = new ConcurrentHashMap(); /* initial seeds for joining the cluster */ private Set seeds_ = new ConcurrentSkipListSet(inetcomparator); @@ -179,7 +179,16 @@ public class Gossiper implements IFailureDetectionEventListener public Set getUnreachableMembers() { - return new HashSet(unreachableEndpoints_); + return unreachableEndpoints_.keySet(); + } + + public long getEndpointDowntime(InetAddress ep) + { + Long downtime = unreachableEndpoints_.get(ep); + if (downtime != null) + return System.currentTimeMillis() - downtime; + else + return 0L; } /** @@ -353,7 +362,7 @@ public class Gossiper implements IFailureDetectionEventListener double prob = unreachableEndpoints / (liveEndpoints + 1); double randDbl = random_.nextDouble(); if ( randDbl < prob ) - sendGossip(message, unreachableEndpoints_); + sendGossip(message, unreachableEndpoints_.keySet()); } } @@ -735,7 +744,7 @@ public class Gossiper implements IFailureDetectionEventListener else { liveEndpoints_.remove(addr); - unreachableEndpoints_.add(addr); + unreachableEndpoints_.put(addr, System.currentTimeMillis()); for (IEndpointStateChangeSubscriber subscriber : subscribers_) subscriber.onDead(addr, epState); } @@ -871,7 +880,7 @@ public class Gossiper implements IFailureDetectionEventListener epState.isAGossiper(true); epState.setHasToken(true); endpointStateMap_.put(ep, epState); - unreachableEndpoints_.add(ep); + unreachableEndpoints_.put(ep, System.currentTimeMillis()); } } diff --git a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java index 7c13068ceb..1a560a02de 100644 --- a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java +++ b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java @@ -25,6 +25,7 @@ import java.util.*; import com.google.common.collect.HashMultimap; import com.google.common.collect.Multimap; +import org.apache.cassandra.gms.Gossiper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -163,6 +164,12 @@ public abstract class AbstractReplicationStrategy { if (map.containsKey(ep)) continue; + if (!StorageProxy.shouldHint(ep)) + { + if (logger.isDebugEnabled()) + logger.debug("not hinting " + ep + " which has been down " + Gossiper.instance.getEndpointDowntime(ep) + "ms"); + continue; + } InetAddress destination = map.isEmpty() ? localAddress diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 6b78762dd9..ec1f435e00 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -75,6 +75,7 @@ public class StorageProxy implements StorageProxyMBean private static final LatencyTracker rangeStats = new LatencyTracker(); private static final LatencyTracker writeStats = new LatencyTracker(); private static boolean hintedHandoffEnabled = DatabaseDescriptor.hintedHandoffEnabled(); + private static int maxHintWindow = DatabaseDescriptor.getMaxHintWindow(); private static final String UNREACHABLE = "UNREACHABLE"; private StorageProxy() {} @@ -182,7 +183,7 @@ public class StorageProxy implements StorageProxyMBean } } responseHandler.addHintCallback(hintedMessage, destination); - + Multimap messages = dcMessages.get(dc); if (messages == null) @@ -190,7 +191,7 @@ public class StorageProxy implements StorageProxyMBean messages = HashMultimap.create(); dcMessages.put(dc, messages); } - + messages.put(hintedMessage, destination); } } @@ -803,6 +804,21 @@ public class StorageProxy implements StorageProxyMBean return hintedHandoffEnabled; } + public int getMaxHintWindow() + { + return maxHintWindow; + } + + public void setMaxHintWindow(int ms) + { + maxHintWindow = ms; + } + + public static boolean shouldHint(InetAddress ep) + { + return Gossiper.instance.getEndpointDowntime(ep) <= maxHintWindow; + } + /** * Performs the truncate operatoin, which effectively deletes all data from * the column family cfname diff --git a/src/java/org/apache/cassandra/service/StorageProxyMBean.java b/src/java/org/apache/cassandra/service/StorageProxyMBean.java index 0c63cf2ba5..8adccec15b 100644 --- a/src/java/org/apache/cassandra/service/StorageProxyMBean.java +++ b/src/java/org/apache/cassandra/service/StorageProxyMBean.java @@ -40,4 +40,6 @@ public interface StorageProxyMBean public boolean getHintedHandoffEnabled(); public void setHintedHandoffEnabled(boolean b); + public int getMaxHintWindow(); + public void setMaxHintWindow(int ms); } From cd9eba9049c0fbbc64669bb79ce93dcd629af9cd Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 20 Jan 2011 23:00:03 +0000 Subject: [PATCH 11/15] set out-of-the-box hint window to one hour git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061560 13f79535-47bb-0310-9956-ffa450edef68 --- conf/cassandra.yaml | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index ca5495d5b9..5f4384341b 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -33,8 +33,7 @@ auto_bootstrap: false hinted_handoff_enabled: true # this defines the maximum amount of time a dead host will have hints # generated. After it has been dead this long, hints will be dropped. -# Maximum is approximately 50 days -max_hint_window_in_ms: 2147483647 +max_hint_window_in_ms: 3600000 # one hour # authentication backend, implementing IAuthenticator; used to identify users authenticator: org.apache.cassandra.auth.AllowAllAuthenticator From e85283501edafc2429a6193271217d36db4b0f91 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 21 Jan 2011 14:51:48 +0000 Subject: [PATCH 12/15] use bytesToHex in debug message git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061830 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/db/commitlog/CommitLog.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLog.java b/src/java/org/apache/cassandra/db/commitlog/CommitLog.java index 761c6e6673..e36e74fbdf 100644 --- a/src/java/org/apache/cassandra/db/commitlog/CommitLog.java +++ b/src/java/org/apache/cassandra/db/commitlog/CommitLog.java @@ -43,6 +43,7 @@ import org.apache.cassandra.db.UnserializableColumnFamilyException; import org.apache.cassandra.io.DeletionService; import org.apache.cassandra.io.util.BufferedRandomAccessFile; import org.apache.cassandra.io.util.FileUtils; +import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.FBUtilities; import org.apache.cassandra.utils.WrappedRunnable; @@ -298,7 +299,7 @@ public class CommitLog if (logger.isDebugEnabled()) logger.debug(String.format("replaying mutation for %s.%s: %s", rm.getTable(), - rm.key(), + ByteBufferUtil.bytesToHex(rm.key()), "{" + StringUtils.join(rm.getColumnFamilies(), ", ") + "}")); final Table table = Table.open(rm.getTable()); tablesRecovered.add(table); From 73d2ad1447e51ecde25f768460e8442b198c4186 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 21 Jan 2011 15:03:54 +0000 Subject: [PATCH 13/15] fix CBL build git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061836 13f79535-47bb-0310-9956-ffa450edef68 --- contrib/bmt_example/CassandraBulkLoader.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/contrib/bmt_example/CassandraBulkLoader.java b/contrib/bmt_example/CassandraBulkLoader.java index e6c072681a..3470d32e37 100644 --- a/contrib/bmt_example/CassandraBulkLoader.java +++ b/contrib/bmt_example/CassandraBulkLoader.java @@ -62,6 +62,7 @@ import java.util.concurrent.TimeoutException; import com.google.common.base.Charsets; import org.apache.cassandra.config.CFMetaData; +import org.apache.cassandra.config.ConfigurationException; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.Column; import org.apache.cassandra.db.ColumnFamily; @@ -112,7 +113,7 @@ public class CassandraBulkLoader { { StorageService.instance.initClient(); } - catch (IOException e) + catch (Exception e) { throw new RuntimeException(e); } From e191a7e6d34f4e63d341504e77f118ade95eafdb Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 21 Jan 2011 16:43:21 +0000 Subject: [PATCH 14/15] more robust error checking on getStorageConfigUrl patch by jbellis git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061894 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/config/DatabaseDescriptor.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 972cd5ea57..1e03c8ca3b 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -101,8 +101,9 @@ public class DatabaseDescriptor try { url = new URL(configUrl); + url.openStream(); // catches well-formed but bogus URLs } - catch (MalformedURLException e) + catch (Exception e) { ClassLoader loader = DatabaseDescriptor.class.getClassLoader(); url = loader.getResource(configUrl); @@ -373,6 +374,7 @@ public class DatabaseDescriptor } catch (UnknownHostException e) { + e.printStackTrace(); logger.error("Fatal error: " + e.getMessage()); System.err.println("Unable to start with unknown hosts configured. Use IP addresses instead of hostnames."); System.exit(2); From a3d120ac1e1a89a08e898de0b41a71b5d9be5cb7 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 21 Jan 2011 16:46:52 +0000 Subject: [PATCH 15/15] r/m printStackTrace git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061897 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/config/DatabaseDescriptor.java | 1 - 1 file changed, 1 deletion(-) diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 1e03c8ca3b..1ea154d436 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -374,7 +374,6 @@ public class DatabaseDescriptor } catch (UnknownHostException e) { - e.printStackTrace(); logger.error("Fatal error: " + e.getMessage()); System.err.println("Unable to start with unknown hosts configured. Use IP addresses instead of hostnames."); System.exit(2);