From e734d837a09a82e43d2e04a53c6921ccd3bc11ac Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 13 Nov 2009 15:57:34 +0000 Subject: [PATCH] fix for when bootstrap source has no data in the range requested patch by jbellis; reviewed by Jaakko Laine for CASSANDRA-541 git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@835891 13f79535-47bb-0310-9956-ffa450edef68 --- .../org/apache/cassandra/io/Streaming.java | 31 ++++++++++++------- 1 file changed, 20 insertions(+), 11 deletions(-) diff --git a/src/java/org/apache/cassandra/io/Streaming.java b/src/java/org/apache/cassandra/io/Streaming.java index c3b00f04fb..cb22649102 100644 --- a/src/java/org/apache/cassandra/io/Streaming.java +++ b/src/java/org/apache/cassandra/io/Streaming.java @@ -69,9 +69,6 @@ public class Streaming private static void transferOneTable(InetAddress target, List sstables, String table) throws IOException { - if (sstables.isEmpty()) - return; - StreamContextManager.StreamContext[] streamContexts = new StreamContextManager.StreamContext[SSTable.FILES_ON_DISK * sstables.size()]; int i = 0; for (SSTableReader sstable : sstables) @@ -91,12 +88,16 @@ public class Streaming if (logger.isDebugEnabled()) logger.debug("Sending a stream initiate message to " + target + " ..."); MessagingService.instance().sendOneWay(message, target); - if (logger.isDebugEnabled()) - logger.debug("Waiting for transfer to " + target + " to complete"); - StreamManager.instance(target).waitForStreamCompletion(); - // reference sstables one more time to make sure it doesn't get GC'd early (causing delete of its files) - if (logger.isDebugEnabled()) - logger.debug("Done with transfer to " + target + " of " + StringUtils.join(sstables, ", ")); + + if (streamContexts.length > 0) + { + if (logger.isDebugEnabled()) + logger.debug("Waiting for transfer to " + target + " to complete"); + StreamManager.instance(target).waitForStreamCompletion(); + // reference sstables one more time to make sure it doesn't get GC'd early (causing delete of its files) + if (logger.isDebugEnabled()) + logger.debug("Done with transfer to " + target + " of " + StringUtils.join(sstables, ", ")); + } } public static class StreamInitiateVerbHandler implements IVerbHandler @@ -119,6 +120,14 @@ public class Streaming StreamInitiateMessage biMsg = StreamInitiateMessage.serializer().deserialize(bufIn); StreamContextManager.StreamContext[] streamContexts = biMsg.getStreamContext(); + if (streamContexts.length == 0 && StorageService.instance().isBootstrapMode()) + { + if (logger.isDebugEnabled()) + logger.debug("no data needed from " + message.getFrom()); + StorageService.instance().removeBootstrapSource(message.getFrom()); + return; + } + Map fileNames = getNewNames(streamContexts); /* * For each of stream context's in the incoming message @@ -142,9 +151,9 @@ public class Streaming Message doneMessage = new Message(FBUtilities.getLocalAddress(), "", StorageService.streamInitiateDoneVerbHandler_, new byte[0] ); MessagingService.instance().sendOneWay(doneMessage, message.getFrom()); } - catch ( IOException ex ) + catch (IOException ex) { - logger.info(LogUtil.throwableToString(ex)); + throw new IOError(ex); } }