From 60b848b43f016091c89c77eb1713ea563136dbf2 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Thu, 20 Jan 2011 22:53:11 +0000 Subject: [PATCH] Add a configurable maximum amount of time to hint for a dead host. Patch by brandonwilliams, reviewed by jbellis for CASSANDRA-1459 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.7@1061557 13f79535-47bb-0310-9956-ffa450edef68 --- conf/cassandra.yaml | 4 ++++ .../org/apache/cassandra/config/Config.java | 1 + .../cassandra/config/DatabaseDescriptor.java | 5 +++++ .../org/apache/cassandra/gms/Gossiper.java | 19 +++++++++++++----- .../locator/AbstractReplicationStrategy.java | 7 +++++++ .../cassandra/service/StorageProxy.java | 20 +++++++++++++++++-- .../cassandra/service/StorageProxyMBean.java | 2 ++ 7 files changed, 51 insertions(+), 7 deletions(-) diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index 8241c78c0d..ca5495d5b9 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -31,6 +31,10 @@ auto_bootstrap: false # See http://wiki.apache.org/cassandra/HintedHandoff hinted_handoff_enabled: true +# this defines the maximum amount of time a dead host will have hints +# generated. After it has been dead this long, hints will be dropped. +# Maximum is approximately 50 days +max_hint_window_in_ms: 2147483647 # authentication backend, implementing IAuthenticator; used to identify users authenticator: org.apache.cassandra.auth.AllowAllAuthenticator diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index def0a5e04b..fd4bca0dbe 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -34,6 +34,7 @@ public class Config public Boolean auto_bootstrap = false; public Boolean hinted_handoff_enabled = true; + public Integer max_hint_window_in_ms = Integer.MAX_VALUE; public String[] seeds; public DiskAccessMode disk_access_mode = DiskAccessMode.auto; diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index f1f23f07ad..972cd5ea57 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -1079,6 +1079,11 @@ public class DatabaseDescriptor return conf.hinted_handoff_enabled; } + public static int getMaxHintWindow() + { + return conf.max_hint_window_in_ms; + } + public static AbstractType getValueValidator(String keyspace, String cf, ByteBuffer column) { return getCFMetaData(keyspace, cf).getValueValidator(column); diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java index de863f831b..312cf85ae5 100644 --- a/src/java/org/apache/cassandra/gms/Gossiper.java +++ b/src/java/org/apache/cassandra/gms/Gossiper.java @@ -128,7 +128,7 @@ public class Gossiper implements IFailureDetectionEventListener private Set liveEndpoints_ = new ConcurrentSkipListSet(inetcomparator); /* unreachable member set */ - private Set unreachableEndpoints_ = new ConcurrentSkipListSet(inetcomparator); + private Map unreachableEndpoints_ = new ConcurrentHashMap(); /* initial seeds for joining the cluster */ private Set seeds_ = new ConcurrentSkipListSet(inetcomparator); @@ -179,7 +179,16 @@ public class Gossiper implements IFailureDetectionEventListener public Set getUnreachableMembers() { - return new HashSet(unreachableEndpoints_); + return unreachableEndpoints_.keySet(); + } + + public long getEndpointDowntime(InetAddress ep) + { + Long downtime = unreachableEndpoints_.get(ep); + if (downtime != null) + return System.currentTimeMillis() - downtime; + else + return 0L; } /** @@ -353,7 +362,7 @@ public class Gossiper implements IFailureDetectionEventListener double prob = unreachableEndpoints / (liveEndpoints + 1); double randDbl = random_.nextDouble(); if ( randDbl < prob ) - sendGossip(message, unreachableEndpoints_); + sendGossip(message, unreachableEndpoints_.keySet()); } } @@ -735,7 +744,7 @@ public class Gossiper implements IFailureDetectionEventListener else { liveEndpoints_.remove(addr); - unreachableEndpoints_.add(addr); + unreachableEndpoints_.put(addr, System.currentTimeMillis()); for (IEndpointStateChangeSubscriber subscriber : subscribers_) subscriber.onDead(addr, epState); } @@ -871,7 +880,7 @@ public class Gossiper implements IFailureDetectionEventListener epState.isAGossiper(true); epState.setHasToken(true); endpointStateMap_.put(ep, epState); - unreachableEndpoints_.add(ep); + unreachableEndpoints_.put(ep, System.currentTimeMillis()); } } diff --git a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java index 7c13068ceb..1a560a02de 100644 --- a/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java +++ b/src/java/org/apache/cassandra/locator/AbstractReplicationStrategy.java @@ -25,6 +25,7 @@ import java.util.*; import com.google.common.collect.HashMultimap; import com.google.common.collect.Multimap; +import org.apache.cassandra.gms.Gossiper; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -163,6 +164,12 @@ public abstract class AbstractReplicationStrategy { if (map.containsKey(ep)) continue; + if (!StorageProxy.shouldHint(ep)) + { + if (logger.isDebugEnabled()) + logger.debug("not hinting " + ep + " which has been down " + Gossiper.instance.getEndpointDowntime(ep) + "ms"); + continue; + } InetAddress destination = map.isEmpty() ? localAddress diff --git a/src/java/org/apache/cassandra/service/StorageProxy.java b/src/java/org/apache/cassandra/service/StorageProxy.java index 6b78762dd9..ec1f435e00 100644 --- a/src/java/org/apache/cassandra/service/StorageProxy.java +++ b/src/java/org/apache/cassandra/service/StorageProxy.java @@ -75,6 +75,7 @@ public class StorageProxy implements StorageProxyMBean private static final LatencyTracker rangeStats = new LatencyTracker(); private static final LatencyTracker writeStats = new LatencyTracker(); private static boolean hintedHandoffEnabled = DatabaseDescriptor.hintedHandoffEnabled(); + private static int maxHintWindow = DatabaseDescriptor.getMaxHintWindow(); private static final String UNREACHABLE = "UNREACHABLE"; private StorageProxy() {} @@ -182,7 +183,7 @@ public class StorageProxy implements StorageProxyMBean } } responseHandler.addHintCallback(hintedMessage, destination); - + Multimap messages = dcMessages.get(dc); if (messages == null) @@ -190,7 +191,7 @@ public class StorageProxy implements StorageProxyMBean messages = HashMultimap.create(); dcMessages.put(dc, messages); } - + messages.put(hintedMessage, destination); } } @@ -803,6 +804,21 @@ public class StorageProxy implements StorageProxyMBean return hintedHandoffEnabled; } + public int getMaxHintWindow() + { + return maxHintWindow; + } + + public void setMaxHintWindow(int ms) + { + maxHintWindow = ms; + } + + public static boolean shouldHint(InetAddress ep) + { + return Gossiper.instance.getEndpointDowntime(ep) <= maxHintWindow; + } + /** * Performs the truncate operatoin, which effectively deletes all data from * the column family cfname diff --git a/src/java/org/apache/cassandra/service/StorageProxyMBean.java b/src/java/org/apache/cassandra/service/StorageProxyMBean.java index 0c63cf2ba5..8adccec15b 100644 --- a/src/java/org/apache/cassandra/service/StorageProxyMBean.java +++ b/src/java/org/apache/cassandra/service/StorageProxyMBean.java @@ -40,4 +40,6 @@ public interface StorageProxyMBean public boolean getHintedHandoffEnabled(); public void setHintedHandoffEnabled(boolean b); + public int getMaxHintWindow(); + public void setMaxHintWindow(int ms); }