From d718dddd6ea1eb4584c4b17151cfdcb11da9236e Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Tue, 4 Jan 2011 21:34:11 +0000 Subject: [PATCH] fix batch mutations post-#1530 patch by tjake; reviewed by jbellis for CASSANDRA-1931 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1055187 13f79535-47bb-0310-9956-ffa450edef68 --- .../cassandra/service/StorageProxy.java | 92 +++++++++++-------- 1 file changed, 56 insertions(+), 36 deletions(-) diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 4cb34e1811..acae142a48 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -123,7 +123,7 @@ public class StorageProxy implements StorageProxyMBean responseHandlers.add(responseHandler); // Multimap that holds onto all the messages and addresses meant for a specific datacenter - Multimap> dcMessages = HashMultimap.create(hintedEndpoints.size(), 10); + Map> dcMessages = new HashMap>(hintedEndpoints.size()); Message unhintedMessage = null; for (Map.Entry> entry : hintedEndpoints.asMap().entrySet()) @@ -150,7 +150,16 @@ public class StorageProxy implements StorageProxyMBean } if (logger.isDebugEnabled()) logger.debug("insert writing key " + FBUtilities.bytesToHex(rm.key()) + " to " + unhintedMessage.getMessageId() + "@" + destination); - dcMessages.put(dc, new Pair(unhintedMessage, destination)); + + + Multimap messages = dcMessages.get(dc); + if (messages == null) + { + messages = HashMultimap.create(); + dcMessages.put(dc, messages); + } + + messages.put(unhintedMessage, destination); } } else @@ -167,7 +176,16 @@ public class StorageProxy implements StorageProxyMBean } } responseHandler.addHintCallback(hintedMessage, destination); - dcMessages.put(dc, new Pair(hintedMessage, destination)); + + Multimap messages = dcMessages.get(dc); + + if (messages == null) + { + messages = HashMultimap.create(); + dcMessages.put(dc, messages); + } + + messages.put(hintedMessage, destination); } } @@ -194,53 +212,55 @@ public class StorageProxy implements StorageProxyMBean /** * for each datacenter, send a message to one node to relay the write to other replicas */ - private static void sendMessages(String localDataCenter, Multimap> dcMessages) + private static void sendMessages(String localDataCenter, Map> dcMessages) throws IOException { - for (Map.Entry>> entry : dcMessages.asMap().entrySet()) + for (Map.Entry> entry: dcMessages.entrySet()) { String dataCenter = entry.getKey(); // Grab a set of all the messages bound for this dataCenter and create an iterator over this set. - Collection> messagesForDataCenter = entry.getValue(); - Iterator> iter = messagesForDataCenter.iterator(); - assert iter.hasNext(); + Map> messagesForDataCenter = entry.getValue().asMap(); - // First endpoint in list is the destination for this group - Pair messageAndDestination = iter.next(); - - Message primaryMessage = messageAndDestination.left; - InetAddress target = messageAndDestination.right; - - // Add all the other destinations that are bound for the same dataCenter as a header in the primary message. - while (iter.hasNext()) + for (Map.Entry> messages: messagesForDataCenter.entrySet()) { - messageAndDestination = iter.next(); - assert messageAndDestination.left == primaryMessage; + Message message = messages.getKey(); + Iterator iter = messages.getValue().iterator(); + assert iter.hasNext(); + + // First endpoint in list is the destination for this group + InetAddress target = iter.next(); + - if (dataCenter.equals(localDataCenter)) + // Add all the other destinations that are bound for the same dataCenter as a header in the primary message. + while (iter.hasNext()) { - // direct write to local DC - assert primaryMessage.getHeader(RowMutation.FORWARD_HEADER) == null; - MessagingService.instance().sendOneWay(primaryMessage, target); - } - else - { - // group all nodes in this DC as forward headers on the primary message - ByteArrayOutputStream bos = new ByteArrayOutputStream(); - DataOutputStream dos = new DataOutputStream(bos); + InetAddress destination = iter.next(); - // append to older addresses - byte[] previousHints = primaryMessage.getHeader(RowMutation.FORWARD_HEADER); - if (previousHints != null) - dos.write(previousHints); + if (dataCenter.equals(localDataCenter)) + { + // direct write to local DC + assert message.getHeader(RowMutation.FORWARD_HEADER) == null; + MessagingService.instance().sendOneWay(message, target); + } + else + { + // group all nodes in this DC as forward headers on the primary message + ByteArrayOutputStream bos = new ByteArrayOutputStream(); + DataOutputStream dos = new DataOutputStream(bos); - dos.write(messageAndDestination.right.getAddress()); - primaryMessage.setHeader(RowMutation.FORWARD_HEADER, bos.toByteArray()); + // append to older addresses + byte[] previousHints = message.getHeader(RowMutation.FORWARD_HEADER); + if (previousHints != null) + dos.write(previousHints); + + dos.write(destination.getAddress()); + message.setHeader(RowMutation.FORWARD_HEADER, bos.toByteArray()); + } } + + MessagingService.instance().sendOneWay(message, target); } - - MessagingService.instance().sendOneWay(primaryMessage, target); } }