From 50b490e046bdc23b23fbb268400abcc41a0de72c Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 12 May 2011 15:31:01 +0000 Subject: [PATCH 01/10] add quote-escaping via backslash to CLI patch by pyaskevich; reviewed by jbellis for CASSANDRA-2623 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1102352 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + src/java/org/apache/cassandra/cli/Cli.g | 51 ++++++++++++++++++- .../cassandra/net/MessagingService.java | 5 +- .../service/AbstractRowResolver.java | 5 +- .../service/DatacenterReadCallback.java | 10 +--- .../cassandra/service/IResponseResolver.java | 3 +- .../service/RangeSliceResponseResolver.java | 9 +--- .../cassandra/service/ReadCallback.java | 23 ++++----- .../org/apache/cassandra/cli/CliTest.java | 5 ++ 9 files changed, 77 insertions(+), 35 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 610e8b473a..a96c77632d 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -23,6 +23,7 @@ * fix counting bloom filter true positives (CASSANDRA-2637) * initialize local ep state prior to gossip startup if needed (CASSANDRA-2638) * fix empty Result with secondary index when limit=1 (CASSANDRA-2628) + * add quote-escaping via backslash to CLI (CASSANDRA-2623) 0.7.5 diff --git a/src/java/org/apache/cassandra/cli/Cli.g b/src/java/org/apache/cassandra/cli/Cli.g index d0b9c99f09..44b6212b4d 100644 --- a/src/java/org/apache/cassandra/cli/Cli.g +++ b/src/java/org/apache/cassandra/cli/Cli.g @@ -581,10 +581,57 @@ Identifier // literals StringLiteral - : - '\'' (~'\'')* '\'' ( '\'' (~'\'')* '\'' )* + : '\'' SingleStringCharacter* '\'' ; +fragment SingleStringCharacter + : ~('\'' | '\\') + | '\\' EscapeSequence + ; + +fragment EscapeSequence + : CharacterEscapeSequence + | '0' + | HexEscapeSequence + | UnicodeEscapeSequence + ; + +fragment CharacterEscapeSequence + : SingleEscapeCharacter + | NonEscapeCharacter + ; + +fragment NonEscapeCharacter + : ~(EscapeCharacter) + ; + +fragment SingleEscapeCharacter + : '\'' | '"' | '\\' | 'b' | 'f' | 'n' | 'r' | 't' | 'v' + ; + +fragment EscapeCharacter + : SingleEscapeCharacter + | DecimalDigit + | 'x' + | 'u' + ; + +fragment HexEscapeSequence + : 'x' HexDigit HexDigit + ; + +fragment UnicodeEscapeSequence + : 'u' HexDigit HexDigit HexDigit HexDigit + ; + +fragment HexDigit + : DecimalDigit | ('a'..'f') | ('A'..'F') + ; + +fragment DecimalDigit + : ('0'..'9') + ; + // // syntactic elements // diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index facda7249b..cc5c40066f 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -259,6 +259,7 @@ public final class MessagingService implements MessagingServiceMBean public String sendRR(Message message, InetAddress to, IMessageCallback cb) { String id = nextId(); + logger_.debug("Message id to {} is {}", to, id); addCallback(cb, id, to); sendOneWay(message, id, to); return id; @@ -266,7 +267,9 @@ public final class MessagingService implements MessagingServiceMBean public void sendOneWay(Message message, InetAddress to) { - sendOneWay(message, nextId(), to); + String id = nextId(); + logger_.debug("Message id to {} is {}", to, id); + sendOneWay(message, id, to); } public void sendReply(Message message, String id, InetAddress to) diff --git a/src/java/org/apache/cassandra/service/AbstractRowResolver.java b/src/java/org/apache/cassandra/service/AbstractRowResolver.java index b13151c5a5..23a22c2c88 100644 --- a/src/java/org/apache/cassandra/service/AbstractRowResolver.java +++ b/src/java/org/apache/cassandra/service/AbstractRowResolver.java @@ -55,16 +55,15 @@ public abstract class AbstractRowResolver implements IResponseResolver this.table = table; } - public void preprocess(Message message) + public ReadResponse preprocess(Message message) { byte[] body = message.getMessageBody(); ByteArrayInputStream bufIn = new ByteArrayInputStream(body); try { ReadResponse result = ReadResponse.serializer().deserialize(new DataInputStream(bufIn)); - if (logger.isDebugEnabled()) - logger.debug("Preprocessed {} response", result.isDigestQuery() ? "digest" : "data"); replies.put(message, result); + return result; } catch (IOException e) { diff --git a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java index 3e6a37ba99..b1b589e31d 100644 --- a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java +++ b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java @@ -48,17 +48,11 @@ public class DatacenterReadCallback extends ReadCallback } @Override - protected boolean waitingFor(Message message) - { - return localdc.equals(snitch.getDatacenter(message.getFrom())); - } - - @Override - protected boolean waitingFor(ReadResponse response) + protected boolean waitingFor(ReadResponse response, InetAddress from) { // cheat and leverage our knowledge that a local read is the only way the ReadResponse // version of this method gets called - return true; + return localdc.equals(snitch.getDatacenter(from)); } @Override diff --git a/src/java/org/apache/cassandra/service/IResponseResolver.java b/src/java/org/apache/cassandra/service/IResponseResolver.java index f4f972a304..0dce952789 100644 --- a/src/java/org/apache/cassandra/service/IResponseResolver.java +++ b/src/java/org/apache/cassandra/service/IResponseResolver.java @@ -20,6 +20,7 @@ package org.apache.cassandra.service; import java.io.IOException; +import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.net.Message; public interface IResponseResolver { @@ -41,6 +42,6 @@ public interface IResponseResolver { */ public T getData() throws IOException; - public void preprocess(Message message); + public ReadResponse preprocess(Message message); public Iterable getMessages(); } diff --git a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java index 09a4d17982..1859889690 100644 --- a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java +++ b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java @@ -24,16 +24,11 @@ import java.util.*; import java.util.concurrent.LinkedBlockingQueue; import com.google.common.collect.AbstractIterator; -import com.google.common.collect.Iterables; -import com.google.common.collect.Iterators; import org.apache.commons.collections.iterators.CollatingIterator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.db.ColumnFamily; -import org.apache.cassandra.db.DecoratedKey; -import org.apache.cassandra.db.RangeSliceReply; -import org.apache.cassandra.db.Row; +import org.apache.cassandra.db.*; import org.apache.cassandra.net.Message; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.ReducingIterator; @@ -114,7 +109,7 @@ public class RangeSliceResponseResolver implements IResponseResolver implements IAsyncCallback public void response(Message message) { - resolver.preprocess(message); - int n = waitingFor(message) + ReadResponse result = resolver.preprocess(message); + int n = waitingFor(result, message.getFrom()) ? received.incrementAndGet() : received.get(); + if (logger.isDebugEnabled()) + logger.debug("{} response; {} qualifying responses seen. Data is {}present", + new Object[] { result.isDigestQuery() ? "digest" : "data", n, resolver.isDataPresent() ? "" : "not " }); if (n >= blockfor && resolver.isDataPresent()) { condition.signal(); @@ -138,19 +142,10 @@ public class ReadCallback implements IAsyncCallback } } - /** - * @return true if the message counts towards the blockfor threshold - * TODO turn the Message into a response so we don't need two versions of this method - */ - protected boolean waitingFor(Message message) - { - return true; - } - /** * @return true if the response counts towards the blockfor threshold */ - protected boolean waitingFor(ReadResponse response) + protected boolean waitingFor(ReadResponse response, InetAddress from) { return true; } @@ -158,7 +153,9 @@ public class ReadCallback implements IAsyncCallback public void response(ReadResponse result) { ((RowDigestResolver) resolver).injectPreProcessed(result); - int n = waitingFor(result) + if (logger.isDebugEnabled()) + logger.debug("Preprocessed {} response", result.isDigestQuery() ? "digest" : "data"); + int n = waitingFor(result, FBUtilities.getLocalAddress()) ? received.incrementAndGet() : received.get(); if (n >= blockfor && resolver.isDataPresent()) diff --git a/test/unit/org/apache/cassandra/cli/CliTest.java b/test/unit/org/apache/cassandra/cli/CliTest.java index bc8450e381..cbd1a9f455 100644 --- a/test/unit/org/apache/cassandra/cli/CliTest.java +++ b/test/unit/org/apache/cassandra/cli/CliTest.java @@ -39,9 +39,14 @@ public class CliTest extends CleanupHelper "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:KEYS, index_name:LongIdxName}];", "set CF1[hello][world] = 123848374878933948398384;", + "set CF1[hello][test_quote] = 'value\\'';", + "set CF1['k\\'ey'][VALUE] = 'VAL';", + "set CF1['k\\'ey'][VALUE] = 'VAL\\'';", "set CF1[hello][-31337] = 'some string value';", "get CF1[hello][-31337];", "get CF1[hello][world];", + "get CF1[hello][test_quote];", + "get CF1['k\\'ey'][VALUE]", "set CF1[hello][-31337] = -23876;", "set CF1[hello][-31337] = long(-23876);", "set CF1[hello][world2] = 15;", From c330ee72e9ea817c21c2424dedddb67d76b9c1b1 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Thu, 12 May 2011 18:18:14 +0000 Subject: [PATCH 02/10] Update files in preparation for 0.7.6 release git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1102406 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 ++ NEWS.txt | 9 +++++++++ build.xml | 2 +- debian/changelog | 6 ++++++ 4 files changed, 18 insertions(+), 1 deletion(-) diff --git a/CHANGES.txt b/CHANGES.txt index a96c77632d..f824c76b4a 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -24,6 +24,8 @@ * initialize local ep state prior to gossip startup if needed (CASSANDRA-2638) * fix empty Result with secondary index when limit=1 (CASSANDRA-2628) * add quote-escaping via backslash to CLI (CASSANDRA-2623) + * fig pig example script (CASSANDRA-2487) + * fix dynamic snitch race in adding latencies (CASSANDRA-2618) 0.7.5 diff --git a/NEWS.txt b/NEWS.txt index 37ef2e0be9..8069b6894c 100644 --- a/NEWS.txt +++ b/NEWS.txt @@ -1,3 +1,12 @@ +0.7.6 +===== + +Upgrading +--------- + - Nothing specific to 0.7.6, but see 0.7.3 Upgrading if upgrading + from earlier than 0.7.1. + + 0.7.5 ===== diff --git a/build.xml b/build.xml index dff5dea944..b3ef9e2e7e 100644 --- a/build.xml +++ b/build.xml @@ -24,7 +24,7 @@ - + diff --git a/debian/changelog b/debian/changelog index 70085dd923..c09d58eec4 100644 --- a/debian/changelog +++ b/debian/changelog @@ -1,3 +1,9 @@ +cassandra (0.7.6) unstable; urgency=low + + * New stable point release. + + -- Sylvain Lebresne Thu, 12 May 2011 20:15:29 +0200 + cassandra (0.7.5) unstable; urgency=low * New stable point release. From c2cf5c4bececf7c2c46c6dfc5ffb9f702835127a Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 12 May 2011 20:33:46 +0000 Subject: [PATCH 03/10] revert work-in-progress accidentally committed git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1102454 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/net/MessagingService.java | 5 +--- .../service/AbstractRowResolver.java | 5 ++-- .../service/DatacenterReadCallback.java | 10 ++++++-- .../cassandra/service/IResponseResolver.java | 3 +-- .../service/RangeSliceResponseResolver.java | 9 ++++++-- .../cassandra/service/ReadCallback.java | 23 +++++++++++-------- 6 files changed, 33 insertions(+), 22 deletions(-) diff --git a/src/java/org/apache/cassandra/net/MessagingService.java b/src/java/org/apache/cassandra/net/MessagingService.java index cc5c40066f..facda7249b 100644 --- a/src/java/org/apache/cassandra/net/MessagingService.java +++ b/src/java/org/apache/cassandra/net/MessagingService.java @@ -259,7 +259,6 @@ public final class MessagingService implements MessagingServiceMBean public String sendRR(Message message, InetAddress to, IMessageCallback cb) { String id = nextId(); - logger_.debug("Message id to {} is {}", to, id); addCallback(cb, id, to); sendOneWay(message, id, to); return id; @@ -267,9 +266,7 @@ public final class MessagingService implements MessagingServiceMBean public void sendOneWay(Message message, InetAddress to) { - String id = nextId(); - logger_.debug("Message id to {} is {}", to, id); - sendOneWay(message, id, to); + sendOneWay(message, nextId(), to); } public void sendReply(Message message, String id, InetAddress to) diff --git a/src/java/org/apache/cassandra/service/AbstractRowResolver.java b/src/java/org/apache/cassandra/service/AbstractRowResolver.java index 23a22c2c88..b13151c5a5 100644 --- a/src/java/org/apache/cassandra/service/AbstractRowResolver.java +++ b/src/java/org/apache/cassandra/service/AbstractRowResolver.java @@ -55,15 +55,16 @@ public abstract class AbstractRowResolver implements IResponseResolver this.table = table; } - public ReadResponse preprocess(Message message) + public void preprocess(Message message) { byte[] body = message.getMessageBody(); ByteArrayInputStream bufIn = new ByteArrayInputStream(body); try { ReadResponse result = ReadResponse.serializer().deserialize(new DataInputStream(bufIn)); + if (logger.isDebugEnabled()) + logger.debug("Preprocessed {} response", result.isDigestQuery() ? "digest" : "data"); replies.put(message, result); - return result; } catch (IOException e) { diff --git a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java index b1b589e31d..3e6a37ba99 100644 --- a/src/java/org/apache/cassandra/service/DatacenterReadCallback.java +++ b/src/java/org/apache/cassandra/service/DatacenterReadCallback.java @@ -48,11 +48,17 @@ public class DatacenterReadCallback extends ReadCallback } @Override - protected boolean waitingFor(ReadResponse response, InetAddress from) + protected boolean waitingFor(Message message) + { + return localdc.equals(snitch.getDatacenter(message.getFrom())); + } + + @Override + protected boolean waitingFor(ReadResponse response) { // cheat and leverage our knowledge that a local read is the only way the ReadResponse // version of this method gets called - return localdc.equals(snitch.getDatacenter(from)); + return true; } @Override diff --git a/src/java/org/apache/cassandra/service/IResponseResolver.java b/src/java/org/apache/cassandra/service/IResponseResolver.java index 0dce952789..f4f972a304 100644 --- a/src/java/org/apache/cassandra/service/IResponseResolver.java +++ b/src/java/org/apache/cassandra/service/IResponseResolver.java @@ -20,7 +20,6 @@ package org.apache.cassandra.service; import java.io.IOException; -import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.net.Message; public interface IResponseResolver { @@ -42,6 +41,6 @@ public interface IResponseResolver { */ public T getData() throws IOException; - public ReadResponse preprocess(Message message); + public void preprocess(Message message); public Iterable getMessages(); } diff --git a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java index 1859889690..09a4d17982 100644 --- a/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java +++ b/src/java/org/apache/cassandra/service/RangeSliceResponseResolver.java @@ -24,11 +24,16 @@ import java.util.*; import java.util.concurrent.LinkedBlockingQueue; import com.google.common.collect.AbstractIterator; +import com.google.common.collect.Iterables; +import com.google.common.collect.Iterators; import org.apache.commons.collections.iterators.CollatingIterator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import org.apache.cassandra.db.*; +import org.apache.cassandra.db.ColumnFamily; +import org.apache.cassandra.db.DecoratedKey; +import org.apache.cassandra.db.RangeSliceReply; +import org.apache.cassandra.db.Row; import org.apache.cassandra.net.Message; import org.apache.cassandra.utils.Pair; import org.apache.cassandra.utils.ReducingIterator; @@ -109,7 +114,7 @@ public class RangeSliceResponseResolver implements IResponseResolver implements IAsyncCallback public void response(Message message) { - ReadResponse result = resolver.preprocess(message); - int n = waitingFor(result, message.getFrom()) + resolver.preprocess(message); + int n = waitingFor(message) ? received.incrementAndGet() : received.get(); - if (logger.isDebugEnabled()) - logger.debug("{} response; {} qualifying responses seen. Data is {}present", - new Object[] { result.isDigestQuery() ? "digest" : "data", n, resolver.isDataPresent() ? "" : "not " }); if (n >= blockfor && resolver.isDataPresent()) { condition.signal(); @@ -142,10 +138,19 @@ public class ReadCallback implements IAsyncCallback } } + /** + * @return true if the message counts towards the blockfor threshold + * TODO turn the Message into a response so we don't need two versions of this method + */ + protected boolean waitingFor(Message message) + { + return true; + } + /** * @return true if the response counts towards the blockfor threshold */ - protected boolean waitingFor(ReadResponse response, InetAddress from) + protected boolean waitingFor(ReadResponse response) { return true; } @@ -153,9 +158,7 @@ public class ReadCallback implements IAsyncCallback public void response(ReadResponse result) { ((RowDigestResolver) resolver).injectPreProcessed(result); - if (logger.isDebugEnabled()) - logger.debug("Preprocessed {} response", result.isDigestQuery() ? "digest" : "data"); - int n = waitingFor(result, FBUtilities.getLocalAddress()) + int n = waitingFor(result) ? received.incrementAndGet() : received.get(); if (n >= blockfor && resolver.isDataPresent()) From 290c6f7d6fddaff6fd80e5537cc961f5e76f42cf Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Thu, 12 May 2011 20:46:14 +0000 Subject: [PATCH 04/10] Start/stop cassandra after more important services such as mdadm in debian packaging. Patch by Paul Cannon, reviewed by brandonwilliams for CASSANDRA-2481 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1102456 13f79535-47bb-0310-9956-ffa450edef68 --- debian/init | 6 ++++-- debian/rules | 2 +- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/debian/init b/debian/init index 709bb6a455..ef3430c0fe 100644 --- a/debian/init +++ b/debian/init @@ -1,8 +1,10 @@ #! /bin/sh ### BEGIN INIT INFO # Provides: cassandra -# Required-Start: $remote_fs -# Required-Stop: $remote_fs +# Required-Start: $remote_fs $network $named $time +# Required-Stop: $remote_fs $network $named $time +# Should-Start: ntp mdadm +# Should-Stop: ntp mdadm # Default-Start: 2 3 4 5 # Default-Stop: 0 1 6 # Short-Description: distributed storage system for structured data diff --git a/debian/rules b/debian/rules index da4cd705b0..9d32d262bd 100755 --- a/debian/rules +++ b/debian/rules @@ -46,7 +46,7 @@ binary-indep: build install dh_testdir dh_testroot dh_installchangelogs - dh_installinit + dh_installinit -u'start 50 2 3 4 5 . stop 50 0 1 6' dh_installdocs README.txt CHANGES.txt NEWS.txt dh_compress dh_fixperms From 81e8f4d3deca0e5672e121e4d1f9c94d3a38233f Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Fri, 13 May 2011 07:58:27 +0000 Subject: [PATCH 05/10] Fix changelog git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1102594 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 ++ 1 file changed, 2 insertions(+) diff --git a/CHANGES.txt b/CHANGES.txt index f824c76b4a..9abc8606a6 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -26,6 +26,8 @@ * add quote-escaping via backslash to CLI (CASSANDRA-2623) * fig pig example script (CASSANDRA-2487) * fix dynamic snitch race in adding latencies (CASSANDRA-2618) + * Start/stop cassandra after more important services such as mdadm in + debian packaging (CASSANDRA-2481) 0.7.5 From 21b5cdb899a7a8ea24b90810e2d32716acc3d7c5 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Mon, 16 May 2011 21:15:12 +0000 Subject: [PATCH 06/10] adjust hinted handoff page size to avoid OOM with large columns patch by jbellis; reviewed by brandonwilliams for CASSANDRA-2652 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1103894 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 7 ++++++- .../cassandra/db/HintedHandOffManager.java | 17 ++++++++++++++--- 2 files changed, 20 insertions(+), 4 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 9abc8606a6..af30314e16 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,3 +1,8 @@ +0.7.7 + * adjust hinted handoff page size to avoid OOM with large columns + (CASSANDRA-2652) + + 0.7.6 * force GC to reclaim disk space on flush, if necessary (CASSANDRA-2404) * move gossip heartbeat back to its own thread (CASSANDRA-2554) @@ -24,7 +29,7 @@ * initialize local ep state prior to gossip startup if needed (CASSANDRA-2638) * fix empty Result with secondary index when limit=1 (CASSANDRA-2628) * add quote-escaping via backslash to CLI (CASSANDRA-2623) - * fig pig example script (CASSANDRA-2487) + * fix pig example script (CASSANDRA-2487) * fix dynamic snitch race in adding latencies (CASSANDRA-2618) * Start/stop cassandra after more important services such as mdadm in debian packaging (CASSANDRA-2481) diff --git a/src/java/org/apache/cassandra/db/HintedHandOffManager.java b/src/java/org/apache/cassandra/db/HintedHandOffManager.java index bbe3cc8f0e..431836c79a 100644 --- a/src/java/org/apache/cassandra/db/HintedHandOffManager.java +++ b/src/java/org/apache/cassandra/db/HintedHandOffManager.java @@ -116,7 +116,7 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean logger_.debug("Created HHOM instance, registered MBean."); } - private static boolean sendMessage(InetAddress endpoint, String tableName, String cfName, ByteBuffer key) throws IOException + private static boolean sendRow(InetAddress endpoint, String tableName, String cfName, ByteBuffer key) throws IOException { if (!Gossiper.instance.isKnownEndpoint(endpoint)) { @@ -131,10 +131,21 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean Table table = Table.open(tableName); DecoratedKey dkey = StorageService.getPartitioner().decorateKey(key); ColumnFamilyStore cfs = table.getColumnFamilyStore(cfName); + + int pageSize = PAGE_SIZE; + // send less columns per page if they are very large + if (cfs.getMeanColumns() > 0) + { + int averageColumnSize = (int) (cfs.getMeanRowSize() / cfs.getMeanColumns()); + pageSize = Math.min(PAGE_SIZE, DatabaseDescriptor.getInMemoryCompactionLimit() / averageColumnSize); + pageSize = Math.max(2, pageSize); // page size of 1 does not allow actual paging b/c of >= behavior on startColumn + logger_.debug("average hinted-row column size is {}; using pageSize of {}", averageColumnSize, pageSize); + } + ByteBuffer startColumn = ByteBufferUtil.EMPTY_BYTE_BUFFER; while (true) { - QueryFilter filter = QueryFilter.getSliceFilter(dkey, new QueryPath(cfs.getColumnFamilyName()), startColumn, ByteBufferUtil.EMPTY_BYTE_BUFFER, false, PAGE_SIZE); + QueryFilter filter = QueryFilter.getSliceFilter(dkey, new QueryPath(cfs.getColumnFamilyName()), startColumn, ByteBufferUtil.EMPTY_BYTE_BUFFER, false, pageSize); ColumnFamily cf = cfs.getColumnFamily(filter); if (pagingFinished(cf, startColumn)) break; @@ -328,7 +339,7 @@ public class HintedHandOffManager implements HintedHandOffManagerMBean for (IColumn tableCF : tableCFs) { String[] parts = getTableAndCFNames(tableCF.name()); - if (sendMessage(endpoint, parts[0], parts[1], keyColumn.name())) + if (sendRow(endpoint, parts[0], parts[1], keyColumn.name())) { deleteHintKey(endpointAsUTF8, keyColumn.name(), tableCF.name(), tableCF.timestamp()); rowsReplayed++; From 2dd17718d250e9c1c53cbd7ae8e6ed8105887dae Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Tue, 17 May 2011 14:53:26 +0000 Subject: [PATCH 07/10] mark BRAF buffer invalid post-flush so we don't re-flush partial buffers again patch by Peter Schuller; reviewed by jbellis for CASSANDRA-2660 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1104305 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 ++ .../io/util/BufferedRandomAccessFile.java | 23 +++++++++++++++---- 2 files changed, 20 insertions(+), 5 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index af30314e16..5628c9c1de 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,6 +1,8 @@ 0.7.7 * adjust hinted handoff page size to avoid OOM with large columns (CASSANDRA-2652) + * mark BRAF buffer invalid post-flush so we don't re-flush partial + buffers again, especially on CL writes (CASSANDRA-2660) 0.7.6 diff --git a/src/java/org/apache/cassandra/io/util/BufferedRandomAccessFile.java b/src/java/org/apache/cassandra/io/util/BufferedRandomAccessFile.java index 00aba8ddc7..439f316b81 100644 --- a/src/java/org/apache/cassandra/io/util/BufferedRandomAccessFile.java +++ b/src/java/org/apache/cassandra/io/util/BufferedRandomAccessFile.java @@ -128,6 +128,9 @@ public class BufferedRandomAccessFile extends RandomAccessFile implements FileDa fd = CLibrary.getfd(this.getFD()); } + /** + * Flush (flush()) whatever writes are pending, and block until the data has been persistently committed (fsync()). + */ public void sync() throws IOException { if (syncNeeded) @@ -150,6 +153,11 @@ public class BufferedRandomAccessFile extends RandomAccessFile implements FileDa } } + /** + * If we are dirty, flush dirty contents to the operating system. Does not imply fsync(). + * + * Currently, for implementation reasons, this also invalidates the buffer. + */ public void flush() throws IOException { if (isDirty) @@ -181,20 +189,25 @@ public class BufferedRandomAccessFile extends RandomAccessFile implements FileDa } + // Remember that we wrote, so we don't write it again on next flush(). + resetBuffer(); + isDirty = false; } } + private void resetBuffer() + { + bufferOffset = current; + validBufferBytes = 0; + } + private void reBuffer() throws IOException { flush(); // synchronizing buffer and file on disk - - bufferOffset = current; + resetBuffer(); if (bufferOffset >= channel.size()) - { - validBufferBytes = 0; return; - } if (bufferOffset < minBufferOffset) minBufferOffset = bufferOffset; From 85ea33912e326cf406d3b9b0c64416b5094839d7 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Wed, 18 May 2011 17:10:14 +0000 Subject: [PATCH 08/10] Fix for dh_installinit syntax for CASSANDRA-2481 Patch by Paul Cannon, reviewed by brandonwilliams git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1124338 13f79535-47bb-0310-9956-ffa450edef68 --- debian/rules | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/debian/rules b/debian/rules index 9d32d262bd..80c8ac2161 100755 --- a/debian/rules +++ b/debian/rules @@ -46,7 +46,7 @@ binary-indep: build install dh_testdir dh_testroot dh_installchangelogs - dh_installinit -u'start 50 2 3 4 5 . stop 50 0 1 6' + dh_installinit -u'start 50 2 3 4 5 . stop 50 0 1 6 .' dh_installdocs README.txt CHANGES.txt NEWS.txt dh_compress dh_fixperms From 4e6c0f6c6bc51ab0e35dbdd0b07017085682bc95 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Thu, 19 May 2011 13:53:06 +0000 Subject: [PATCH 09/10] Allow limiting row slices in the cli patch by slebresne; reviewed by xedin for CASSANDRA-2646 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1124780 13f79535-47bb-0310-9956-ffa450edef68 --- src/java/org/apache/cassandra/cli/Cli.g | 8 ++-- .../org/apache/cassandra/cli/CliClient.java | 45 +++++++++++++++---- 2 files changed, 40 insertions(+), 13 deletions(-) diff --git a/src/java/org/apache/cassandra/cli/Cli.g b/src/java/org/apache/cassandra/cli/Cli.g index 44b6212b4d..e954838c69 100644 --- a/src/java/org/apache/cassandra/cli/Cli.g +++ b/src/java/org/apache/cassandra/cli/Cli.g @@ -228,10 +228,10 @@ exitStatement ; getStatement - : GET columnFamilyExpr ('AS' typeIdentifier)? - -> ^(NODE_THRIFT_GET columnFamilyExpr ( ^(CONVERT_TO_TYPE typeIdentifier) )? ) - | GET columnFamily 'WHERE' getCondition ('AND' getCondition)* ('LIMIT' limit=IntegerLiteral)* - -> ^(NODE_THRIFT_GET_WITH_CONDITIONS columnFamily ^(CONDITIONS getCondition+) ^(NODE_LIMIT $limit)*) + : GET columnFamilyExpr ('AS' typeIdentifier)? ('LIMIT' limit=IntegerLiteral)? + -> ^(NODE_THRIFT_GET columnFamilyExpr ( ^(CONVERT_TO_TYPE typeIdentifier) )? ^(NODE_LIMIT $limit)?) + | GET columnFamily 'WHERE' getCondition ('AND' getCondition)* ('LIMIT' limit=IntegerLiteral)? + -> ^(NODE_THRIFT_GET_WITH_CONDITIONS columnFamily ^(CONDITIONS getCondition+) ^(NODE_LIMIT $limit)?) ; getCondition diff --git a/src/java/org/apache/cassandra/cli/CliClient.java b/src/java/org/apache/cassandra/cli/CliClient.java index dbaff6da3d..2611e08e50 100644 --- a/src/java/org/apache/cassandra/cli/CliClient.java +++ b/src/java/org/apache/cassandra/cli/CliClient.java @@ -306,7 +306,7 @@ public class CliClient extends CliUserHelp sessionState.out.println(String.format("%s removed.", (columnSpecCnt == 0) ? "row" : "column")); } - private void doSlice(String keyspace, ByteBuffer key, String columnFamily, byte[] superColumnName) + private void doSlice(String keyspace, ByteBuffer key, String columnFamily, byte[] superColumnName, int limit) throws InvalidRequestException, UnavailableException, TimedOutException, TException, IllegalAccessException, NotFoundException, InstantiationException, NoSuchFieldException { @@ -314,7 +314,7 @@ public class CliClient extends CliUserHelp if(superColumnName != null) parent.setSuper_column(superColumnName); - SliceRange range = new SliceRange(ByteBufferUtil.EMPTY_BYTE_BUFFER, ByteBufferUtil.EMPTY_BYTE_BUFFER, false, 1000000); + SliceRange range = new SliceRange(ByteBufferUtil.EMPTY_BYTE_BUFFER, ByteBufferUtil.EMPTY_BYTE_BUFFER, false, limit); List columns = thriftClient.get_slice(key, parent, new SlicePredicate().setColumn_names(null).setSlice_range(range), consistencyLevel); AbstractType validator; @@ -401,10 +401,39 @@ public class CliClient extends CliUserHelp byte[] superColumnName = null; ByteBuffer columnName; + Tree typeTree = null; + Tree limitTree = null; + + int limit = 1000000; + + if (statement.getChildCount() >= 2) + { + if (statement.getChild(1).getType() == CliParser.CONVERT_TO_TYPE) + { + typeTree = statement.getChild(1).getChild(0); + if (statement.getChildCount() == 3) + limitTree = statement.getChild(2).getChild(0); + } + else + { + limitTree = statement.getChild(1).getChild(0); + } + } + + if (limitTree != null) + { + limit = Integer.parseInt(limitTree.getText()); + + if (limit == 0) + { + throw new IllegalArgumentException("LIMIT should be greater than zero."); + } + } + // table.cf['key'] -- row slice if (columnSpecCnt == 0) { - doSlice(keySpace, key, columnFamily, superColumnName); + doSlice(keySpace, key, columnFamily, superColumnName, limit); return; } // table.cf['key']['column'] -- slice of a super, or get of a standard @@ -415,7 +444,7 @@ public class CliClient extends CliUserHelp if (isSuper) { superColumnName = columnName.array(); - doSlice(keySpace, key, columnFamily, superColumnName); + doSlice(keySpace, key, columnFamily, superColumnName, limit); return; } } @@ -433,7 +462,7 @@ public class CliClient extends CliUserHelp } AbstractType validator = getValidatorForValue(cfDef, TBaseHelper.byteBufferToByteArray(columnName)); - + // Perform a get() ColumnPath path = new ColumnPath(columnFamily); if(superColumnName != null) path.setSuper_column(superColumnName); @@ -451,14 +480,12 @@ public class CliClient extends CliUserHelp byte[] columnValue = column.getValue(); String valueAsString; - + // we have ^(CONVERT_TO_TYPE ) inside of GET statement // which means that we should try to represent byte[] value according // to specified type - if (statement.getChildCount() == 2) + if (typeTree != null) { - // getting ^(CONVERT_TO_TYPE ) tree - Tree typeTree = statement.getChild(1).getChild(0); // .getText() will give us String typeName = CliUtils.unescapeSQLString(typeTree.getText()); // building AbstractType from From 050d129179b99fa88d25a65afed78d0fedb24f1e Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Thu, 19 May 2011 17:10:49 +0000 Subject: [PATCH 10/10] don't perform HH to client-mode patch by jbellis; reviewed by brandonwilliams for CASSANDRA-2668 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1125002 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 1 + src/java/org/apache/cassandra/gms/EndpointState.java | 2 +- src/java/org/apache/cassandra/gms/Gossiper.java | 2 +- src/java/org/apache/cassandra/service/StorageService.java | 2 +- 4 files changed, 4 insertions(+), 3 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 5628c9c1de..501d98fe68 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -3,6 +3,7 @@ (CASSANDRA-2652) * mark BRAF buffer invalid post-flush so we don't re-flush partial buffers again, especially on CL writes (CASSANDRA-2660) + * don't perform HH to client-mode [storageproxy] nodes (CASSANDRA-2668) 0.7.6 diff --git a/src/java/org/apache/cassandra/gms/EndpointState.java b/src/java/org/apache/cassandra/gms/EndpointState.java index aa17917301..14ab23f843 100644 --- a/src/java/org/apache/cassandra/gms/EndpointState.java +++ b/src/java/org/apache/cassandra/gms/EndpointState.java @@ -136,7 +136,7 @@ public class EndpointState hasToken_ = value; } - public boolean getHasToken() + public boolean hasToken() { return hasToken_; } diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index 057b560564..1c17fe686b 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -431,7 +431,7 @@ public class Gossiper implements IFailureDetectionEventListener // check if this is a fat client. fat clients are removed automatically from // gosip after FatClientTimeout - if (!epState.getHasToken() && !epState.isAlive() && (duration > FatClientTimeout_)) + if (!epState.hasToken() && !epState.isAlive() && (duration > FatClientTimeout_)) { if (StorageService.instance.getTokenMetadata().isMember(endpoint)) epState.setHasToken(true); diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 0cb8ad63bb..07768db9b4 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -1130,7 +1130,7 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe public void onAlive(InetAddress endpoint, EndpointState state) { - if (!isClientMode) + if (!isClientMode && state.hasToken()) deliverHints(endpoint); }