From 8e72ac4f9fbb11cd6aab12f2d262743dcb394f49 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Mon, 20 Apr 2009 20:05:07 +0000 Subject: [PATCH] rename get_cf -> readColumnFamily; ReadMessage -> ReadCommand. [Message message = ReadMessage.readMessage(readMessage) is just plain confusing] patch by jbellis; reviewed by Eric Evans for #88 git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@766838 13f79535-47bb-0310-9956-ffa450edef68 --- .../db/{ReadMessage.java => ReadCommand.java} | 41 +++--- .../apache/cassandra/db/ReadVerbHandler.java | 49 +++---- .../cassandra/service/CassandraServer.java | 7 +- .../cassandra/service/ConsistencyManager.java | 30 ++-- .../service/MultiQuorumResponseHandler.java | 14 +- .../cassandra/service/StorageProxy.java | 130 +++++++++--------- .../apache/cassandra/test/DataImporter.java | 6 +- src/org/apache/cassandra/test/StressTest.java | 17 +-- .../apache/cassandra/db/ReadMessageTest.java | 10 +- 9 files changed, 142 insertions(+), 162 deletions(-) rename src/org/apache/cassandra/db/{ReadMessage.java => ReadCommand.java} (78%) diff --git a/src/org/apache/cassandra/db/ReadMessage.java b/src/org/apache/cassandra/db/ReadCommand.java similarity index 78% rename from src/org/apache/cassandra/db/ReadMessage.java rename to src/org/apache/cassandra/db/ReadCommand.java index 8fa2803a28..85460bef42 100644 --- a/src/org/apache/cassandra/db/ReadMessage.java +++ b/src/org/apache/cassandra/db/ReadCommand.java @@ -28,7 +28,6 @@ import java.util.List; import org.apache.commons.lang.StringUtils; -import org.apache.cassandra.continuations.Suspendable; import org.apache.cassandra.io.ICompactSerializer; import org.apache.cassandra.net.Message; import org.apache.cassandra.service.StorageService; @@ -38,26 +37,26 @@ import org.apache.cassandra.service.StorageService; * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) */ -public class ReadMessage implements Serializable +public class ReadCommand implements Serializable { - private static ICompactSerializer serializer_; + private static ICompactSerializer serializer_; public static final String doRepair_ = "READ-REPAIR"; - + static { - serializer_ = new ReadMessageSerializer(); + serializer_ = new ReadCommandSerializer(); } - static ICompactSerializer serializer() + static ICompactSerializer serializer() { return serializer_; } - public static Message makeReadMessage(ReadMessage readMessage) throws IOException + public static Message makeReadMessage(ReadCommand readCommand) throws IOException { ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream( bos ); - ReadMessage.serializer().serialize(readMessage, dos); + ReadCommand.serializer().serialize(readCommand, dos); Message message = new Message(StorageService.getLocalStorageEndPoint(), StorageService.readStage_, StorageService.readVerbHandler_, new Object[]{bos.toByteArray()}); return message; } @@ -71,24 +70,24 @@ public class ReadMessage implements Serializable private List columns_ = new ArrayList(); private boolean isDigestQuery_ = false; - private ReadMessage() + private ReadCommand() { } - public ReadMessage(String table, String key) + public ReadCommand(String table, String key) { table_ = table; key_ = key; } - public ReadMessage(String table, String key, String columnFamily_column) + public ReadCommand(String table, String key, String columnFamily_column) { table_ = table; key_ = key; columnFamily_column_ = columnFamily_column; } - public ReadMessage(String table, String key, String columnFamily, List columns) + public ReadCommand(String table, String key, String columnFamily, List columns) { table_ = table; key_ = key; @@ -96,7 +95,7 @@ public class ReadMessage implements Serializable columns_ = columns; } - public ReadMessage(String table, String key, String columnFamily_column, int start, int count) + public ReadCommand(String table, String key, String columnFamily_column, int start, int count) { table_ = table; key_ = key; @@ -105,7 +104,7 @@ public class ReadMessage implements Serializable count_ = count; } - public ReadMessage(String table, String key, String columnFamily_column, long sinceTimestamp) + public ReadCommand(String table, String key, String columnFamily_column, long sinceTimestamp) { table_ = table; key_ = key; @@ -173,9 +172,9 @@ public class ReadMessage implements Serializable } } -class ReadMessageSerializer implements ICompactSerializer +class ReadCommandSerializer implements ICompactSerializer { - public void serialize(ReadMessage rm, DataOutputStream dos) throws IOException + public void serialize(ReadCommand rm, DataOutputStream dos) throws IOException { dos.writeUTF(rm.table()); dos.writeUTF(rm.key()); @@ -196,7 +195,7 @@ class ReadMessageSerializer implements ICompactSerializer } } - public ReadMessage deserialize(DataInputStream dis) throws IOException + public ReadCommand deserialize(DataInputStream dis) throws IOException { String table = dis.readUTF(); String key = dis.readUTF(); @@ -214,18 +213,18 @@ class ReadMessageSerializer implements ICompactSerializer dis.readFully(bytes); columns.add( new String(bytes) ); } - ReadMessage rm = null; + ReadCommand rm = null; if ( columns.size() > 0 ) { - rm = new ReadMessage(table, key, columnFamily_column, columns); + rm = new ReadCommand(table, key, columnFamily_column, columns); } else if( sinceTimestamp > 0 ) { - rm = new ReadMessage(table, key, columnFamily_column, sinceTimestamp); + rm = new ReadCommand(table, key, columnFamily_column, sinceTimestamp); } else { - rm = new ReadMessage(table, key, columnFamily_column, start, count); + rm = new ReadCommand(table, key, columnFamily_column, start, count); } rm.setIsDigestQuery(isDigest); return rm; diff --git a/src/org/apache/cassandra/db/ReadVerbHandler.java b/src/org/apache/cassandra/db/ReadVerbHandler.java index 4bfabf5f68..e4b39bdd1e 100644 --- a/src/org/apache/cassandra/db/ReadVerbHandler.java +++ b/src/org/apache/cassandra/db/ReadVerbHandler.java @@ -19,12 +19,9 @@ package org.apache.cassandra.db; import java.io.IOException; -import java.nio.ByteBuffer; -import java.util.Collection; import java.util.List; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.continuations.Suspendable; import org.apache.cassandra.io.DataInputBuffer; import org.apache.cassandra.io.DataOutputBuffer; import org.apache.cassandra.net.EndPoint; @@ -34,8 +31,6 @@ import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.service.StorageService; import org.apache.cassandra.utils.LogUtil; import org.apache.log4j.Logger; -import org.apache.cassandra.net.*; -import org.apache.cassandra.utils.*; /** * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) @@ -77,30 +72,30 @@ public class ReadVerbHandler implements IVerbHandler try { - ReadMessage readMessage = ReadMessage.serializer().deserialize(readCtx.bufIn_); - Table table = Table.open(readMessage.table()); + ReadCommand readCommand = ReadCommand.serializer().deserialize(readCtx.bufIn_); + Table table = Table.open(readCommand.table()); Row row = null; long start = System.currentTimeMillis(); - if( readMessage.columnFamily_column() == null ) - row = table.get(readMessage.key()); + if( readCommand.columnFamily_column() == null ) + row = table.get(readCommand.key()); else { - if(readMessage.getColumnNames().size() == 0) + if(readCommand.getColumnNames().size() == 0) { - if(readMessage.count() > 0 && readMessage.start() >= 0) - row = table.getRow(readMessage.key(), readMessage.columnFamily_column(), readMessage.start(), readMessage.count()); + if(readCommand.count() > 0 && readCommand.start() >= 0) + row = table.getRow(readCommand.key(), readCommand.columnFamily_column(), readCommand.start(), readCommand.count()); else - row = table.getRow(readMessage.key(), readMessage.columnFamily_column()); + row = table.getRow(readCommand.key(), readCommand.columnFamily_column()); } else { - row = table.getRow(readMessage.key(), readMessage.columnFamily_column(), readMessage.getColumnNames()); + row = table.getRow(readCommand.key(), readCommand.columnFamily_column(), readCommand.getColumnNames()); } } logger_.info("getRow() TIME: " + (System.currentTimeMillis() - start) + " ms."); start = System.currentTimeMillis(); ReadResponseMessage readResponseMessage = null; - if(readMessage.isDigestQuery()) + if(readCommand.isDigestQuery()) { readResponseMessage = new ReadResponseMessage(table.getTableName(), row.digest()); } @@ -108,7 +103,7 @@ public class ReadVerbHandler implements IVerbHandler { readResponseMessage = new ReadResponseMessage(table.getTableName(), row); } - readResponseMessage.setIsDigestQuery(readMessage.isDigestQuery()); + readResponseMessage.setIsDigestQuery(readCommand.isDigestQuery()); /* serialize the ReadResponseMessage. */ readCtx.bufOut_.reset(); @@ -126,9 +121,9 @@ public class ReadVerbHandler implements IVerbHandler logger_.info("ReadVerbHandler TIME 2: " + (System.currentTimeMillis() - start) + " ms."); /* Do read repair if header of the message says so */ - String repair = new String( message.getHeader(ReadMessage.doRepair_) ); - if ( repair.equals( ReadMessage.doRepair_ ) ) - doReadRepair(row, readMessage); + String repair = new String( message.getHeader(ReadCommand.doRepair_) ); + if ( repair.equals( ReadCommand.doRepair_ ) ) + doReadRepair(row, readCommand); } catch ( IOException ex) { @@ -140,29 +135,29 @@ public class ReadVerbHandler implements IVerbHandler } } - private void doReadRepair(Row row, ReadMessage readMessage) + private void doReadRepair(Row row, ReadCommand readCommand) { if ( DatabaseDescriptor.getConsistencyCheck() ) { - List endpoints = StorageService.instance().getNLiveStorageEndPoint(readMessage.key()); + List endpoints = StorageService.instance().getNLiveStorageEndPoint(readCommand.key()); /* Remove the local storage endpoint from the list. */ endpoints.remove( StorageService.getLocalStorageEndPoint() ); - if(readMessage.getColumnNames().size() == 0) + if(readCommand.getColumnNames().size() == 0) { - if( readMessage.start() >= 0 && readMessage.count() < Integer.MAX_VALUE) + if( readCommand.start() >= 0 && readCommand.count() < Integer.MAX_VALUE) { - StorageService.instance().doConsistencyCheck(row, endpoints, readMessage.columnFamily_column(), readMessage.start(), readMessage.count()); + StorageService.instance().doConsistencyCheck(row, endpoints, readCommand.columnFamily_column(), readCommand.start(), readCommand.count()); } - if( readMessage.sinceTimestamp() > 0) + if( readCommand.sinceTimestamp() > 0) { - StorageService.instance().doConsistencyCheck(row, endpoints, readMessage.columnFamily_column(), readMessage.sinceTimestamp()); + StorageService.instance().doConsistencyCheck(row, endpoints, readCommand.columnFamily_column(), readCommand.sinceTimestamp()); } } else { - StorageService.instance().doConsistencyCheck(row, endpoints, readMessage.columnFamily_column(), readMessage.getColumnNames()); + StorageService.instance().doConsistencyCheck(row, endpoints, readCommand.columnFamily_column(), readCommand.getColumnNames()); } } } diff --git a/src/org/apache/cassandra/service/CassandraServer.java b/src/org/apache/cassandra/service/CassandraServer.java index 3ba334bb64..050581ca57 100644 --- a/src/org/apache/cassandra/service/CassandraServer.java +++ b/src/org/apache/cassandra/service/CassandraServer.java @@ -37,7 +37,6 @@ import org.apache.cassandra.db.ColumnFamily; import org.apache.cassandra.db.IColumn; import org.apache.cassandra.db.Row; import org.apache.cassandra.db.RowMutation; -import org.apache.cassandra.db.Column; import org.apache.cassandra.utils.LogUtil; import org.apache.thrift.TException; @@ -81,7 +80,7 @@ public class CassandraServer implements Cassandra.Iface } } - protected ColumnFamily get_cf(String tablename, String key, String columnFamily, List columNames) throws CassandraException, TException + protected ColumnFamily readColumnFamily(String tablename, String key, String columnFamily, List columNames) throws CassandraException, TException { ColumnFamily cfamily = null; try @@ -205,7 +204,7 @@ public class CassandraServer implements Cassandra.Iface try { validateTable(tablename); - ColumnFamily cfamily = get_cf(tablename, key, columnFamily, columnNames); + ColumnFamily cfamily = readColumnFamily(tablename, key, columnFamily, columnNames); if (cfamily == null) { logger_.info("ERROR ColumnFamily " + columnFamily + " is missing.....: " @@ -486,7 +485,7 @@ public class CassandraServer implements Cassandra.Iface try { validateTable(tablename); - ColumnFamily cfamily = get_cf(tablename, key, columnFamily, superColumnNames); + ColumnFamily cfamily = readColumnFamily(tablename, key, columnFamily, superColumnNames); if (cfamily == null) { logger_.info("ERROR ColumnFamily " + columnFamily + " is missing.....: "+" key:" + key diff --git a/src/org/apache/cassandra/service/ConsistencyManager.java b/src/org/apache/cassandra/service/ConsistencyManager.java index f8a87df28b..33238aaa8c 100644 --- a/src/org/apache/cassandra/service/ConsistencyManager.java +++ b/src/org/apache/cassandra/service/ConsistencyManager.java @@ -22,13 +22,9 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Arrays; import java.util.List; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; -import org.apache.cassandra.concurrent.DebuggableScheduledThreadPoolExecutor; -import org.apache.cassandra.concurrent.ThreadFactoryImpl; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.ReadMessage; +import org.apache.cassandra.db.ReadCommand; import org.apache.cassandra.db.ReadResponseMessage; import org.apache.cassandra.db.Row; import org.apache.cassandra.io.DataInputBuffer; @@ -96,9 +92,9 @@ class ConsistencyManager implements Runnable replicas_.add(StorageService.getLocalStorageEndPoint()); IAsyncCallback responseHandler = new DataRepairHandler(ConsistencyManager.this.replicas_.size(), readResponseResolver); String table = DatabaseDescriptor.getTables().get(0); - ReadMessage readMessage = constructReadMessage(false); + ReadCommand readCommand = constructReadMessage(false); // ReadMessage readMessage = new ReadMessage(table, row_.key(), columnFamily_); - Message message = ReadMessage.makeReadMessage(readMessage); + Message message = ReadCommand.makeReadMessage(readCommand); MessagingService.getMessagingInstance().sendRR(message, replicas_.toArray( new EndPoint[0] ), responseHandler); } } @@ -187,10 +183,10 @@ class ConsistencyManager implements Runnable public void run() { logger_.debug(" Run the consistency checks for " + columnFamily_); - ReadMessage readMessageDigestOnly = constructReadMessage(true); + ReadCommand readCommandDigestOnly = constructReadMessage(true); try { - Message messageDigestOnly = ReadMessage.makeReadMessage(readMessageDigestOnly); + Message messageDigestOnly = ReadCommand.makeReadMessage(readCommandDigestOnly); IAsyncCallback digestResponseHandler = new DigestResponseHandler(); MessagingService.getMessagingInstance().sendRR(messageDigestOnly, replicas_.toArray(new EndPoint[0]), digestResponseHandler); } @@ -200,32 +196,32 @@ class ConsistencyManager implements Runnable } } - private ReadMessage constructReadMessage(boolean isDigestQuery) + private ReadCommand constructReadMessage(boolean isDigestQuery) { - ReadMessage readMessage = null; + ReadCommand readCommand = null; String table = DatabaseDescriptor.getTables().get(0); if(columnNames_.size() == 0) { if( start_ >= 0 && count_ < Integer.MAX_VALUE) { - readMessage = new ReadMessage(table, row_.key(), columnFamily_, start_, count_); + readCommand = new ReadCommand(table, row_.key(), columnFamily_, start_, count_); } else if(sinceTimestamp_ > 0) { - readMessage = new ReadMessage(table, row_.key(), columnFamily_, sinceTimestamp_); + readCommand = new ReadCommand(table, row_.key(), columnFamily_, sinceTimestamp_); } else { - readMessage = new ReadMessage(table, row_.key(), columnFamily_); + readCommand = new ReadCommand(table, row_.key(), columnFamily_); } } else { - readMessage = new ReadMessage(table, row_.key(), columnFamily_, columnNames_); + readCommand = new ReadCommand(table, row_.key(), columnFamily_, columnNames_); } - readMessage.setIsDigestQuery(isDigestQuery); - return readMessage; + readCommand.setIsDigestQuery(isDigestQuery); + return readCommand; } } diff --git a/src/org/apache/cassandra/service/MultiQuorumResponseHandler.java b/src/org/apache/cassandra/service/MultiQuorumResponseHandler.java index f3331d1144..2d4dfe3a5a 100644 --- a/src/org/apache/cassandra/service/MultiQuorumResponseHandler.java +++ b/src/org/apache/cassandra/service/MultiQuorumResponseHandler.java @@ -28,7 +28,7 @@ import java.util.concurrent.locks.*; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.ReadMessage; +import org.apache.cassandra.db.ReadCommand; import org.apache.cassandra.db.Row; import org.apache.cassandra.net.EndPoint; import org.apache.cassandra.net.IAsyncCallback; @@ -47,7 +47,7 @@ public class MultiQuorumResponseHandler implements IAsyncCallback private Lock lock_ = new ReentrantLock(); private Condition condition_; /* This maps the keys to the original data read messages */ - private Map readMessages_ = new HashMap(); + private Map readMessages_ = new HashMap(); /* This maps the key to its set of replicas */ private Map endpoints_ = new HashMap(); /* This maps the groupId to the individual callback for the set of messages */ @@ -129,10 +129,10 @@ public class MultiQuorumResponseHandler implements IAsyncCallback { if ( DatabaseDescriptor.getConsistencyCheck()) { - ReadMessage readMessage = readMessages_.get(key); - readMessage.setIsDigestQuery(false); - Message messageRepair = ReadMessage.makeReadMessage(readMessage); - EndPoint[] endpoints = MultiQuorumResponseHandler.this.endpoints_.get( readMessage.key() ); + ReadCommand readCommand = readMessages_.get(key); + readCommand.setIsDigestQuery(false); + Message messageRepair = ReadCommand.makeReadMessage(readCommand); + EndPoint[] endpoints = MultiQuorumResponseHandler.this.endpoints_.get( readCommand.key() ); Message[][] messages = new Message[][]{ {messageRepair, messageRepair, messageRepair} }; EndPoint[][] epList = new EndPoint[][]{ endpoints }; MessagingService.getMessagingInstance().sendRR(messages, epList, MultiQuorumResponseHandler.this); @@ -140,7 +140,7 @@ public class MultiQuorumResponseHandler implements IAsyncCallback } } - public MultiQuorumResponseHandler(Map readMessages, Map endpoints) + public MultiQuorumResponseHandler(Map readMessages, Map endpoints) { condition_ = lock_.newCondition(); readMessages_ = readMessages; diff --git a/src/org/apache/cassandra/service/StorageProxy.java b/src/org/apache/cassandra/service/StorageProxy.java index 5b189be7e9..c3d0b66cb5 100644 --- a/src/org/apache/cassandra/service/StorageProxy.java +++ b/src/org/apache/cassandra/service/StorageProxy.java @@ -30,16 +30,14 @@ import java.util.concurrent.TimeoutException; import org.apache.commons.lang.StringUtils; import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.ReadMessage; +import org.apache.cassandra.db.ReadCommand; import org.apache.cassandra.db.ReadResponseMessage; import org.apache.cassandra.db.Row; import org.apache.cassandra.db.RowMutation; -import org.apache.cassandra.db.RowMutationMessage; import org.apache.cassandra.db.Table; import org.apache.cassandra.db.TouchMessage; import org.apache.cassandra.io.DataInputBuffer; import org.apache.cassandra.net.EndPoint; -import org.apache.cassandra.net.IAsyncCallback; import org.apache.cassandra.net.IAsyncResult; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; @@ -145,13 +143,13 @@ public class StorageProxy } } - private static Map constructMessages(Map readMessages) throws IOException + private static Map constructMessages(Map readMessages) throws IOException { Map messages = new HashMap(); Set keys = readMessages.keySet(); for ( String key : keys ) { - Message message = ReadMessage.makeReadMessage( readMessages.get(key) ); + Message message = ReadCommand.makeReadMessage( readMessages.get(key) ); messages.put(key, message); } return messages; @@ -182,7 +180,7 @@ public class StorageProxy * @throws IOException * @throws TimeoutException */ - public static Map doReadProtocol(Map readMessages) throws IOException,TimeoutException + public static Map doReadProtocol(Map readMessages) throws IOException,TimeoutException { Map rows = new HashMap(); Set keys = readMessages.keySet(); @@ -206,14 +204,14 @@ public class StorageProxy return rows; } - public static Row doReadProtocol(String key, ReadMessage readMessage) throws IOException,TimeoutException + public static Row doReadProtocol(String key, ReadCommand readCommand) throws IOException,TimeoutException { Row row = null; EndPoint endPoint = StorageService.instance().findSuitableEndPoint(key); if(endPoint != null) { - Message message = ReadMessage.makeReadMessage(readMessage); - message.addHeader(ReadMessage.doRepair_, ReadMessage.doRepair_.getBytes()); + Message message = ReadCommand.makeReadMessage(readCommand); + message.addHeader(ReadCommand.doRepair_, ReadCommand.doRepair_.getBytes()); IAsyncResult iar = MessagingService.getMessagingInstance().sendRR(message, endPoint); Object[] result = iar.get(DatabaseDescriptor.getRpcTimeout(), TimeUnit.MILLISECONDS); byte[] body = (byte[])result[0]; @@ -310,9 +308,9 @@ public class StorageProxy } if(!foundLocal && consistencyLevel == StorageService.ConsistencyLevel.WEAK) { - ReadMessage readMessage = null; - readMessage = new ReadMessage(tablename, key, columnFamily, columnNames); - return doReadProtocol(key, readMessage); + ReadCommand readCommand = null; + readCommand = new ReadCommand(tablename, key, columnFamily, columnNames); + return doReadProtocol(key, readCommand); } else { @@ -349,9 +347,9 @@ public class StorageProxy } if(!foundLocal && consistencyLevel == StorageService.ConsistencyLevel.WEAK) { - ReadMessage readMessage = null; - readMessage = new ReadMessage(tablename, key, columnFamily, start, count); - return doReadProtocol(key, readMessage); + ReadCommand readCommand = null; + readCommand = new ReadCommand(tablename, key, columnFamily, start, count); + return doReadProtocol(key, readCommand); } else { @@ -408,9 +406,9 @@ public class StorageProxy } if(!foundLocal && consistencyLevel == StorageService.ConsistencyLevel.WEAK) { - ReadMessage readMessage = null; - readMessage = new ReadMessage(tablename, key, columnFamily, sinceTimestamp); - return doReadProtocol(key, readMessage); + ReadCommand readCommand = null; + readCommand = new ReadCommand(tablename, key, columnFamily, sinceTimestamp); + return doReadProtocol(key, readCommand); } else { @@ -436,12 +434,12 @@ public class StorageProxy { long startTime = System.currentTimeMillis(); // TODO: throw a thrift exception if we do not have N nodes - ReadMessage readMessage = new ReadMessage(tablename, key, columnFamily, columns); + ReadCommand readCommand = new ReadCommand(tablename, key, columnFamily, columns); - ReadMessage readMessageDigestOnly = new ReadMessage(tablename, key, columnFamily, columns); - readMessageDigestOnly.setIsDigestQuery(true); + ReadCommand readCommandDigestOnly = new ReadCommand(tablename, key, columnFamily, columns); + readCommandDigestOnly.setIsDigestQuery(true); - Row row = StorageProxy.doStrongReadProtocol(key, readMessage, readMessageDigestOnly); + Row row = StorageProxy.doStrongReadProtocol(key, readCommand, readCommandDigestOnly); logger_.debug("readProtocol: " + (System.currentTimeMillis() - startTime) + " ms."); return row; } @@ -466,27 +464,27 @@ public class StorageProxy { long startTime = System.currentTimeMillis(); // TODO: throw a thrift exception if we do not have N nodes - ReadMessage readMessage = null; - ReadMessage readMessageDigestOnly = null; + ReadCommand readCommand = null; + ReadCommand readCommandDigestOnly = null; if( start >= 0 && count < Integer.MAX_VALUE) { - readMessage = new ReadMessage(tablename, key, columnFamily, start, count); + readCommand = new ReadCommand(tablename, key, columnFamily, start, count); } else { - readMessage = new ReadMessage(tablename, key, columnFamily); + readCommand = new ReadCommand(tablename, key, columnFamily); } - Message message = ReadMessage.makeReadMessage(readMessage); + Message message = ReadCommand.makeReadMessage(readCommand); if( start >= 0 && count < Integer.MAX_VALUE) { - readMessageDigestOnly = new ReadMessage(tablename, key, columnFamily, start, count); + readCommandDigestOnly = new ReadCommand(tablename, key, columnFamily, start, count); } else { - readMessageDigestOnly = new ReadMessage(tablename, key, columnFamily); + readCommandDigestOnly = new ReadCommand(tablename, key, columnFamily); } - readMessageDigestOnly.setIsDigestQuery(true); - Row row = doStrongReadProtocol(key, readMessage, readMessageDigestOnly); + readCommandDigestOnly.setIsDigestQuery(true); + Row row = doStrongReadProtocol(key, readCommand, readCommandDigestOnly); logger_.debug("readProtocol: " + (System.currentTimeMillis() - startTime) + " ms."); return row; } @@ -507,27 +505,27 @@ public class StorageProxy Map rows = new HashMap(); long startTime = System.currentTimeMillis(); // TODO: throw a thrift exception if we do not have N nodes - Map readMessages = new HashMap(); + Map readMessages = new HashMap(); for (String key : keys ) { - ReadMessage[] readMessage = new ReadMessage[2]; + ReadCommand[] readCommand = new ReadCommand[2]; if( start >= 0 && count < Integer.MAX_VALUE) { - readMessage[0] = new ReadMessage(tablename, key, columnFamily, start, count); + readCommand[0] = new ReadCommand(tablename, key, columnFamily, start, count); } else { - readMessage[0] = new ReadMessage(tablename, key, columnFamily); + readCommand[0] = new ReadCommand(tablename, key, columnFamily); } if( start >= 0 && count < Integer.MAX_VALUE) { - readMessage[1] = new ReadMessage(tablename, key, columnFamily, start, count); + readCommand[1] = new ReadCommand(tablename, key, columnFamily, start, count); } else { - readMessage[1] = new ReadMessage(tablename, key, columnFamily); + readCommand[1] = new ReadCommand(tablename, key, columnFamily); } - readMessage[1].setIsDigestQuery(true); + readCommand[1].setIsDigestQuery(true); } rows = doStrongReadProtocol(readMessages); logger_.debug("readProtocol: " + (System.currentTimeMillis() - startTime) + " ms."); @@ -538,13 +536,13 @@ public class StorageProxy { long startTime = System.currentTimeMillis(); // TODO: throw a thrift exception if we do not have N nodes - ReadMessage readMessage = null; - ReadMessage readMessageDigestOnly = null; - readMessage = new ReadMessage(tablename, key, columnFamily, sinceTimestamp); - Message message = ReadMessage.makeReadMessage(readMessage); - readMessageDigestOnly = new ReadMessage(tablename, key, columnFamily, sinceTimestamp); - readMessageDigestOnly.setIsDigestQuery(true); - Row row = doStrongReadProtocol(key, readMessage, readMessageDigestOnly); + ReadCommand readCommand = null; + ReadCommand readCommandDigestOnly = null; + readCommand = new ReadCommand(tablename, key, columnFamily, sinceTimestamp); + Message message = ReadCommand.makeReadMessage(readCommand); + readCommandDigestOnly = new ReadCommand(tablename, key, columnFamily, sinceTimestamp); + readCommandDigestOnly.setIsDigestQuery(true); + Row row = doStrongReadProtocol(key, readCommand, readCommandDigestOnly); logger_.debug("readProtocol: " + (System.currentTimeMillis() - startTime) + " ms."); return row; } @@ -555,11 +553,11 @@ public class StorageProxy * param @ readMessage - the read message to get the actual data * param @ readMessageDigest - the read message to get the digest. */ - private static Row doStrongReadProtocol(String key, ReadMessage readMessage, ReadMessage readMessageDigest) throws IOException, TimeoutException + private static Row doStrongReadProtocol(String key, ReadCommand readCommand, ReadCommand readCommandDigest) throws IOException, TimeoutException { Row row = null; - Message message = ReadMessage.makeReadMessage(readMessage); - Message messageDigestOnly = ReadMessage.makeReadMessage(readMessageDigest); + Message message = ReadCommand.makeReadMessage(readCommand); + Message messageDigestOnly = ReadCommand.makeReadMessage(readCommandDigest); IResponseResolver readResponseResolver = new ReadResponseResolver(); QuorumResponseHandler quorumResponseHandler = new QuorumResponseHandler( @@ -607,9 +605,9 @@ public class StorageProxy QuorumResponseHandler quorumResponseHandlerRepair = new QuorumResponseHandler( DatabaseDescriptor.getReplicationFactor(), readResponseResolverRepair); - readMessage.setIsDigestQuery(false); + readCommand.setIsDigestQuery(false); logger_.info("DigestMismatchException: " + key); - Message messageRepair = ReadMessage.makeReadMessage(readMessage); + Message messageRepair = ReadCommand.makeReadMessage(readCommand); MessagingService.getMessagingInstance().sendRR(messageRepair, endPoints, quorumResponseHandlerRepair); try { @@ -628,7 +626,7 @@ public class StorageProxy return row; } - private static Map constructReplicaMessages(Map readMessages) throws IOException + private static Map constructReplicaMessages(Map readMessages) throws IOException { Map messages = new HashMap(); Set keys = readMessages.keySet(); @@ -636,21 +634,21 @@ public class StorageProxy for ( String key : keys ) { Message[] msg = new Message[DatabaseDescriptor.getReplicationFactor()]; - ReadMessage[] readMessage = readMessages.get(key); - msg[0] = ReadMessage.makeReadMessage( readMessage[0] ); + ReadCommand[] readCommand = readMessages.get(key); + msg[0] = ReadCommand.makeReadMessage( readCommand[0] ); for ( int i = 1; i < msg.length; ++i ) { - msg[i] = ReadMessage.makeReadMessage( readMessage[1] ); + msg[i] = ReadCommand.makeReadMessage( readCommand[1] ); } } return messages; } - private static MultiQuorumResponseHandler dispatchMessages(Map readMessages, Map messages) throws IOException + private static MultiQuorumResponseHandler dispatchMessages(Map readMessages, Map messages) throws IOException { Set keys = messages.keySet(); /* This maps the keys to the original data read messages */ - Map readMessage = new HashMap(); + Map readMessage = new HashMap(); /* This maps the keys to their respective endpoints/replicas */ Map endpoints = new HashMap(); /* Groups the messages that need to be sent to the individual keys */ @@ -700,7 +698,7 @@ public class StorageProxy * @return map containing key ---> Row * @throws IOException, TimeoutException */ - private static Map doStrongReadProtocol(Map readMessages) throws IOException + private static Map doStrongReadProtocol(Map readMessages) throws IOException { Map rows = new HashMap(); /* Construct the messages to be sent to the replicas */ @@ -769,11 +767,11 @@ public class StorageProxy { Row row = null; long startTime = System.currentTimeMillis(); - Map readMessages = new HashMap(); + Map readMessages = new HashMap(); for ( String key : keys ) { - ReadMessage readMessage = new ReadMessage(tablename, key, columnFamily, columns); - readMessages.put(key, readMessage); + ReadCommand readCommand = new ReadCommand(tablename, key, columnFamily, columns); + readMessages.put(key, readCommand); } /* Performs the multiget in parallel */ Map rows = doReadProtocol(readMessages); @@ -850,11 +848,11 @@ public class StorageProxy { Row row = null; long startTime = System.currentTimeMillis(); - Map readMessages = new HashMap(); + Map readMessages = new HashMap(); for ( String key : keys ) { - ReadMessage readMessage = new ReadMessage(tablename, key, columnFamily, start, count); - readMessages.put(key, readMessage); + ReadCommand readCommand = new ReadCommand(tablename, key, columnFamily, start, count); + readMessages.put(key, readCommand); } /* Performs the multiget in parallel */ Map rows = doReadProtocol(readMessages); @@ -919,11 +917,11 @@ public class StorageProxy { Row row = null; long startTime = System.currentTimeMillis(); - Map readMessages = new HashMap(); + Map readMessages = new HashMap(); for ( String key : keys ) { - ReadMessage readMessage = new ReadMessage(tablename, key, columnFamily, sinceTimestamp); - readMessages.put(key, readMessage); + ReadCommand readCommand = new ReadCommand(tablename, key, columnFamily, sinceTimestamp); + readMessages.put(key, readCommand); } /* Performs the multiget in parallel */ Map rows = doReadProtocol(readMessages); diff --git a/src/org/apache/cassandra/test/DataImporter.java b/src/org/apache/cassandra/test/DataImporter.java index 9bb8585087..0caf99bde4 100644 --- a/src/org/apache/cassandra/test/DataImporter.java +++ b/src/org/apache/cassandra/test/DataImporter.java @@ -43,7 +43,7 @@ import org.apache.cassandra.concurrent.DebuggableScheduledThreadPoolExecutor; import org.apache.cassandra.concurrent.ThreadFactoryImpl; import org.apache.cassandra.db.ColumnFamily; import org.apache.cassandra.db.IColumn; -import org.apache.cassandra.db.ReadMessage; +import org.apache.cassandra.db.ReadCommand; import org.apache.cassandra.db.ReadResponseMessage; import org.apache.cassandra.db.Row; import org.apache.cassandra.db.RowMutation; @@ -880,10 +880,10 @@ public class DataImporter { key = user + ":1"; } - ReadMessage readMessage = new ReadMessage(tablename_, key); + ReadCommand readCommand = new ReadCommand(tablename_, key); Message message = new Message(from_, StorageService.readStage_, StorageService.readVerbHandler_, - new Object[] { readMessage }); + new Object[] {readCommand}); IAsyncResult iar = MessagingService.getMessagingInstance().sendRR( message, to_); Object[] result = iar.get(); diff --git a/src/org/apache/cassandra/test/StressTest.java b/src/org/apache/cassandra/test/StressTest.java index 9c999dcc6e..d6c88a76b8 100644 --- a/src/org/apache/cassandra/test/StressTest.java +++ b/src/org/apache/cassandra/test/StressTest.java @@ -19,10 +19,7 @@ package org.apache.cassandra.test; import java.io.IOException; -import java.text.DateFormat; -import java.text.SimpleDateFormat; import java.util.ArrayList; -import java.util.Date; import java.util.HashMap; import java.util.List; import java.util.Random; @@ -30,17 +27,13 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; -import org.apache.cassandra.analytics.AnalyticsContext; import org.apache.cassandra.concurrent.DebuggableThreadPoolExecutor; import org.apache.cassandra.concurrent.ThreadFactoryImpl; -import org.apache.cassandra.config.DatabaseDescriptor; -import org.apache.cassandra.db.Memtable; -import org.apache.cassandra.db.ReadMessage; +import org.apache.cassandra.db.ReadCommand; import org.apache.cassandra.db.Row; import org.apache.cassandra.db.RowMutation; import org.apache.cassandra.db.RowMutationMessage; import org.apache.cassandra.net.EndPoint; -import org.apache.cassandra.net.IAsyncResult; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.service.Cassandra; @@ -147,7 +140,7 @@ public class StressTest } - public void readLoad(ReadMessage readMessage) + public void readLoad(ReadCommand readCommand) { IResponseResolver readResponseResolver = new ReadResponseResolver(); QuorumResponseHandler quorumResponseHandler = new QuorumResponseHandler( @@ -155,7 +148,7 @@ public class StressTest readResponseResolver); Message message = new Message(from_, StorageService.readStage_, StorageService.readVerbHandler_, - new Object[] { readMessage }); + new Object[] {readCommand}); MessagingService.getMessagingInstance().sendOneWay(message, to_); /*IAsyncResult iar = MessagingService.getMessagingInstance().sendRR(message, to_); try @@ -187,7 +180,7 @@ public class StressTest String stringKey = new Integer(key).toString(); stringKey = stringKey + keyFix_ ; int j = random.nextInt(columns) + 1; - ReadMessage rm = new ReadMessage(tablename_, stringKey, columnFamilyColumn_ + ":" + columnFix_ + j); + ReadCommand rm = new ReadCommand(tablename_, stringKey, columnFamilyColumn_ + ":" + columnFix_ + j); readLoad(rm); if ( requestsPerSecond_ > 1000) Thread.sleep(0, 1000000000/requestsPerSecond_); @@ -257,7 +250,7 @@ public class StressTest stringKey = stringKey + keyFix_ ; int i = random.nextInt(superColumns) + 1; int j = random.nextInt(columns) + 1; - ReadMessage rm = new ReadMessage(tablename_, stringKey, columnFamilySuperColumn_ + ":" + superColumnFix_ + i + ":" + columnFix_ + j); + ReadCommand rm = new ReadCommand(tablename_, stringKey, columnFamilySuperColumn_ + ":" + superColumnFix_ + i + ":" + columnFix_ + j); readLoad(rm); } } diff --git a/test/org/apache/cassandra/db/ReadMessageTest.java b/test/org/apache/cassandra/db/ReadMessageTest.java index 583e795585..f517c9a6dd 100644 --- a/test/org/apache/cassandra/db/ReadMessageTest.java +++ b/test/org/apache/cassandra/db/ReadMessageTest.java @@ -16,16 +16,16 @@ public class ReadMessageTest colList.add("col1"); colList.add("col2"); - ReadMessage rm = new ReadMessage("Table1", "row1", "foo", colList); - ReadMessage rm2 = serializeAndDeserializeReadMessage(rm); + ReadCommand rm = new ReadCommand("Table1", "row1", "foo", colList); + ReadCommand rm2 = serializeAndDeserializeReadMessage(rm); assert rm2.toString().equals(rm.toString()); } - private ReadMessage serializeAndDeserializeReadMessage(ReadMessage rm) + private ReadCommand serializeAndDeserializeReadMessage(ReadCommand rm) { - ReadMessage rm2 = null; - ReadMessageSerializer rms = (ReadMessageSerializer) ReadMessage.serializer(); + ReadCommand rm2 = null; + ReadCommandSerializer rms = (ReadCommandSerializer) ReadCommand.serializer(); DataOutputBuffer dos = new DataOutputBuffer(); DataInputBuffer dis = new DataInputBuffer();