diff --git a/src/org/apache/cassandra/db/ReadResponseMessage.java b/src/org/apache/cassandra/db/ReadResponse.java similarity index 74% rename from src/org/apache/cassandra/db/ReadResponseMessage.java rename to src/org/apache/cassandra/db/ReadResponse.java index 8efe8a9af2..ab8c96a528 100644 --- a/src/org/apache/cassandra/db/ReadResponseMessage.java +++ b/src/org/apache/cassandra/db/ReadResponse.java @@ -35,25 +35,25 @@ import org.apache.cassandra.service.StorageService; * The table name is needed so that we can use it to create repairs. * Author : Avinash Lakshman ( alakshman@facebook.com) & Prashant Malik ( pmalik@facebook.com ) */ -public class ReadResponseMessage implements Serializable +public class ReadResponse implements Serializable { -private static ICompactSerializer serializer_; - +private static ICompactSerializer serializer_; + static { - serializer_ = new ReadResponseMessageSerializer(); + serializer_ = new ReadResponseSerializer(); } - public static ICompactSerializer serializer() + public static ICompactSerializer serializer() { return serializer_; } - public static Message makeReadResponseMessage(ReadResponseMessage readResponseMessage) throws IOException + public static Message makeReadResponseMessage(ReadResponse readResponse) throws IOException { ByteArrayOutputStream bos = new ByteArrayOutputStream(); DataOutputStream dos = new DataOutputStream( bos ); - ReadResponseMessage.serializer().serialize(readResponseMessage, dos); + ReadResponse.serializer().serialize(readResponse, dos); Message message = new Message(StorageService.getLocalStorageEndPoint(), MessagingService.responseStage_, MessagingService.responseVerbHandler_, new Object[]{bos.toByteArray()}); return message; } @@ -63,13 +63,13 @@ private static ICompactSerializer serializer_; private byte[] digest_ = new byte[0]; private boolean isDigestQuery_ = false; - public ReadResponseMessage(String table, byte[] digest ) + public ReadResponse(String table, byte[] digest ) { table_ = table; digest_= digest; } - public ReadResponseMessage(String table, Row row) + public ReadResponse(String table, Row row) { table_ = table; row_ = row; @@ -101,9 +101,9 @@ private static ICompactSerializer serializer_; } } -class ReadResponseMessageSerializer implements ICompactSerializer +class ReadResponseSerializer implements ICompactSerializer { - public void serialize(ReadResponseMessage rm, DataOutputStream dos) throws IOException + public void serialize(ReadResponse rm, DataOutputStream dos) throws IOException { dos.writeUTF(rm.table()); dos.writeInt(rm.digest().length); @@ -116,7 +116,7 @@ class ReadResponseMessageSerializer implements ICompactSerializer try { long start = System.currentTimeMillis(); - ReadResponseMessage result = ReadResponseMessage.serializer().deserialize(bufIn); + ReadResponse result = ReadResponse.serializer().deserialize(bufIn); logger_.debug( "Response deserialization time : " + (System.currentTimeMillis() - start) + " ms."); if(!result.isDigestQuery()) { @@ -168,7 +168,7 @@ public class ReadResponseResolver implements IResponseResolver bufIn.reset(body, body.length); try { - ReadResponseMessage result = ReadResponseMessage.serializer().deserialize(bufIn); + ReadResponse result = ReadResponse.serializer().deserialize(bufIn); if(!result.isDigestQuery()) { isDataPresent = true; diff --git a/src/org/apache/cassandra/service/StorageProxy.java b/src/org/apache/cassandra/service/StorageProxy.java index bea90e5d08..e5a13606e6 100644 --- a/src/org/apache/cassandra/service/StorageProxy.java +++ b/src/org/apache/cassandra/service/StorageProxy.java @@ -31,7 +31,7 @@ import org.apache.commons.lang.StringUtils; import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.db.ReadCommand; -import org.apache.cassandra.db.ReadResponseMessage; +import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.db.Row; import org.apache.cassandra.db.RowMutation; import org.apache.cassandra.db.Table; @@ -197,8 +197,8 @@ public class StorageProxy byte[] body = (byte[])result[0]; DataInputBuffer bufIn = new DataInputBuffer(); bufIn.reset(body, body.length); - ReadResponseMessage responseMessage = ReadResponseMessage.serializer().deserialize(bufIn); - Row row = responseMessage.row(); + ReadResponse response = ReadResponse.serializer().deserialize(bufIn); + Row row = response.row(); rows.put(row.key(), row); } return rows; @@ -217,8 +217,8 @@ public class StorageProxy byte[] body = (byte[])result[0]; DataInputBuffer bufIn = new DataInputBuffer(); bufIn.reset(body, body.length); - ReadResponseMessage responseMessage = ReadResponseMessage.serializer().deserialize(bufIn); - row = responseMessage.row(); + ReadResponse response = ReadResponse.serializer().deserialize(bufIn); + row = response.row(); } else { diff --git a/src/org/apache/cassandra/test/DataImporter.java b/src/org/apache/cassandra/test/DataImporter.java index 0caf99bde4..7b72ce214f 100644 --- a/src/org/apache/cassandra/test/DataImporter.java +++ b/src/org/apache/cassandra/test/DataImporter.java @@ -44,7 +44,7 @@ import org.apache.cassandra.concurrent.ThreadFactoryImpl; import org.apache.cassandra.db.ColumnFamily; import org.apache.cassandra.db.IColumn; import org.apache.cassandra.db.ReadCommand; -import org.apache.cassandra.db.ReadResponseMessage; +import org.apache.cassandra.db.ReadResponse; import org.apache.cassandra.db.Row; import org.apache.cassandra.db.RowMutation; import org.apache.cassandra.db.RowMutationMessage; @@ -887,8 +887,8 @@ public class DataImporter { IAsyncResult iar = MessagingService.getMessagingInstance().sendRR( message, to_); Object[] result = iar.get(); - ReadResponseMessage readResponseMessage = (ReadResponseMessage) result[0]; - Row row = readResponseMessage.row(); + ReadResponse readResponse = (ReadResponse) result[0]; + Row row = readResponse.row(); if (row == null) { logger_.debug("ERROR No row for this key .....: " + line); Thread.sleep(1000/requestsPerSecond_, 1000%requestsPerSecond_);