mirror of https://github.com/apache/cassandra
Changes in here to enable multiget() support.
git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@759025 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
parent
f98301a68d
commit
f5d1a1289c
|
|
@ -25,8 +25,6 @@ import java.io.IOException;
|
|||
import java.io.Serializable;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import javax.xml.bind.annotation.XmlElement;
|
||||
|
||||
import org.apache.cassandra.continuations.Suspendable;
|
||||
import org.apache.cassandra.io.ICompactSerializer;
|
||||
import org.apache.cassandra.net.Message;
|
||||
|
|
@ -40,6 +38,7 @@ import org.apache.cassandra.service.StorageService;
|
|||
public class ReadMessage implements Serializable
|
||||
{
|
||||
private static ICompactSerializer<ReadMessage> serializer_;
|
||||
public static final String doRepair_ = "READ-REPAIR";
|
||||
|
||||
static
|
||||
{
|
||||
|
|
@ -60,28 +59,13 @@ public class ReadMessage implements Serializable
|
|||
return message;
|
||||
}
|
||||
|
||||
@XmlElement(name="Table")
|
||||
private String table_;
|
||||
|
||||
@XmlElement(name="Key")
|
||||
private String key_;
|
||||
|
||||
@XmlElement(name="ColumnFamily")
|
||||
private String columnFamily_column_ = null;
|
||||
|
||||
@XmlElement(name="start")
|
||||
private int start_ = -1;
|
||||
|
||||
@XmlElement(name="count")
|
||||
private int count_ = -1 ;
|
||||
|
||||
@XmlElement(name="sinceTimestamp")
|
||||
private long sinceTimestamp_ = -1 ;
|
||||
|
||||
@XmlElement(name="columnNames")
|
||||
private List<String> columns_ = new ArrayList<String>();
|
||||
|
||||
@XmlElement(name="isDigestQuery")
|
||||
private boolean isDigestQuery_ = false;
|
||||
|
||||
private ReadMessage()
|
||||
|
|
@ -131,7 +115,7 @@ public class ReadMessage implements Serializable
|
|||
return table_;
|
||||
}
|
||||
|
||||
String key()
|
||||
public String key()
|
||||
{
|
||||
return key_;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -23,9 +23,6 @@ import java.io.DataInputStream;
|
|||
import java.io.DataOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.Serializable;
|
||||
|
||||
import javax.xml.bind.annotation.XmlElement;
|
||||
|
||||
import org.apache.cassandra.io.ICompactSerializer;
|
||||
import org.apache.cassandra.net.Message;
|
||||
import org.apache.cassandra.net.MessagingService;
|
||||
|
|
@ -61,32 +58,25 @@ private static ICompactSerializer<ReadResponseMessage> serializer_;
|
|||
return message;
|
||||
}
|
||||
|
||||
@XmlElement(name = "Table")
|
||||
private String table_;
|
||||
|
||||
@XmlElement(name = "Row")
|
||||
private Row row_;
|
||||
|
||||
@XmlElement(name = "Digest")
|
||||
private byte[] digest_ = new byte[0];
|
||||
|
||||
@XmlElement(name="isDigestQuery")
|
||||
private boolean isDigestQuery_ = false;
|
||||
|
||||
private ReadResponseMessage() {
|
||||
}
|
||||
|
||||
public ReadResponseMessage(String table, byte[] digest ) {
|
||||
public ReadResponseMessage(String table, byte[] digest )
|
||||
{
|
||||
table_ = table;
|
||||
digest_= digest;
|
||||
}
|
||||
|
||||
public ReadResponseMessage(String table, Row row) {
|
||||
public ReadResponseMessage(String table, Row row)
|
||||
{
|
||||
table_ = table;
|
||||
row_ = row;
|
||||
}
|
||||
|
||||
public String table() {
|
||||
public String table()
|
||||
{
|
||||
return table_;
|
||||
}
|
||||
|
||||
|
|
@ -95,7 +85,8 @@ private static ICompactSerializer<ReadResponseMessage> serializer_;
|
|||
return row_;
|
||||
}
|
||||
|
||||
public byte[] digest() {
|
||||
public byte[] digest()
|
||||
{
|
||||
return digest_;
|
||||
}
|
||||
|
||||
|
|
@ -110,7 +101,6 @@ private static ICompactSerializer<ReadResponseMessage> serializer_;
|
|||
}
|
||||
}
|
||||
|
||||
|
||||
class ReadResponseMessageSerializer implements ICompactSerializer<ReadResponseMessage>
|
||||
{
|
||||
public void serialize(ReadResponseMessage rm, DataOutputStream dos) throws IOException
|
||||
|
|
|
|||
|
|
@ -21,10 +21,13 @@ 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;
|
||||
import org.apache.cassandra.net.IVerbHandler;
|
||||
import org.apache.cassandra.net.Message;
|
||||
import org.apache.cassandra.net.MessagingService;
|
||||
|
|
@ -40,7 +43,7 @@ import org.apache.cassandra.utils.*;
|
|||
|
||||
public class ReadVerbHandler implements IVerbHandler
|
||||
{
|
||||
private static class ReadContext
|
||||
protected static class ReadContext
|
||||
{
|
||||
protected DataInputBuffer bufIn_ = new DataInputBuffer();
|
||||
protected DataOutputBuffer bufOut_ = new DataOutputBuffer();
|
||||
|
|
@ -48,7 +51,17 @@ public class ReadVerbHandler implements IVerbHandler
|
|||
|
||||
private static Logger logger_ = Logger.getLogger( ReadVerbHandler.class );
|
||||
/* We use this so that we can reuse the same row mutation context for the mutation. */
|
||||
private static ThreadLocal<ReadContext> tls_ = new InheritableThreadLocal<ReadContext>();
|
||||
private static ThreadLocal<ReadVerbHandler.ReadContext> tls_ = new InheritableThreadLocal<ReadVerbHandler.ReadContext>();
|
||||
|
||||
protected static ReadVerbHandler.ReadContext getCurrentReadContext()
|
||||
{
|
||||
return tls_.get();
|
||||
}
|
||||
|
||||
protected static void setCurrentReadContext(ReadVerbHandler.ReadContext readContext)
|
||||
{
|
||||
tls_.set(readContext);
|
||||
}
|
||||
|
||||
public void doVerb(Message message)
|
||||
{
|
||||
|
|
@ -90,7 +103,6 @@ public class ReadVerbHandler implements IVerbHandler
|
|||
if(readMessage.isDigestQuery())
|
||||
{
|
||||
readResponseMessage = new ReadResponseMessage(table.getTableName(), row.digest());
|
||||
|
||||
}
|
||||
else
|
||||
{
|
||||
|
|
@ -111,8 +123,12 @@ public class ReadVerbHandler implements IVerbHandler
|
|||
|
||||
Message response = message.getReply( StorageService.getLocalStorageEndPoint(), new Object[]{bytes} );
|
||||
MessagingService.getMessagingInstance().sendOneWay(response, message.getFrom());
|
||||
logger_.info("ReadVerbHandler TIME 2: " + (System.currentTimeMillis() - start)
|
||||
+ " ms.");
|
||||
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);
|
||||
}
|
||||
catch ( IOException ex)
|
||||
{
|
||||
|
|
@ -123,4 +139,31 @@ public class ReadVerbHandler implements IVerbHandler
|
|||
logger_.info( LogUtil.throwableToString(ex) );
|
||||
}
|
||||
}
|
||||
|
||||
private void doReadRepair(Row row, ReadMessage readMessage)
|
||||
{
|
||||
if ( DatabaseDescriptor.getConsistencyCheck() )
|
||||
{
|
||||
List<EndPoint> endpoints = StorageService.instance().getNLiveStorageEndPoint(readMessage.key());
|
||||
/* Remove the local storage endpoint from the list. */
|
||||
endpoints.remove( StorageService.getLocalStorageEndPoint() );
|
||||
|
||||
if(readMessage.getColumnNames().size() == 0)
|
||||
{
|
||||
if( readMessage.start() >= 0 && readMessage.count() < Integer.MAX_VALUE)
|
||||
{
|
||||
StorageService.instance().doConsistencyCheck(row, endpoints, readMessage.columnFamily_column(), readMessage.start(), readMessage.count());
|
||||
}
|
||||
|
||||
if( readMessage.sinceTimestamp() > 0)
|
||||
{
|
||||
StorageService.instance().doConsistencyCheck(row, endpoints, readMessage.columnFamily_column(), readMessage.sinceTimestamp());
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
StorageService.instance().doConsistencyCheck(row, endpoints, readMessage.columnFamily_column(), readMessage.getColumnNames());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue