Author: jdcryans
Date: Thu Sep 15 21:48:59 2011
New Revision: 1171286

URL: http://svn.apache.org/viewvc?rev=1171286&view=rev
Log:
   HBASE-4363  [replication] ReplicationSource won't close if failing
               to contact the sink (JD and Lars Hofhansl)

Modified:
    hbase/trunk/CHANGES.txt
    
hbase/trunk/src/main/java/org/apache/hadoop/hbase/replication/regionserver/ReplicationSource.java

Modified: hbase/trunk/CHANGES.txt
URL: 
http://svn.apache.org/viewvc/hbase/trunk/CHANGES.txt?rev=1171286&r1=1171285&r2=1171286&view=diff
==============================================================================
--- hbase/trunk/CHANGES.txt (original)
+++ hbase/trunk/CHANGES.txt Thu Sep 15 21:48:59 2011
@@ -266,6 +266,8 @@ Release 0.91.0 - Unreleased
    HBASE-4351  If from Admin we try to unassign a region forcefully,
                though a valid region name is given the master is not able
                to identify the region to unassign (Ramkrishna)
+   HBASE-4363  [replication] ReplicationSource won't close if failing
+               to contact the sink (JD and Lars Hofhansl)
 
   IMPROVEMENTS
    HBASE-3290  Max Compaction Size (Nicolas Spiegelberg via Stack)  

Modified: 
hbase/trunk/src/main/java/org/apache/hadoop/hbase/replication/regionserver/ReplicationSource.java
URL: 
http://svn.apache.org/viewvc/hbase/trunk/src/main/java/org/apache/hadoop/hbase/replication/regionserver/ReplicationSource.java?rev=1171286&r1=1171285&r2=1171286&view=diff
==============================================================================
--- 
hbase/trunk/src/main/java/org/apache/hadoop/hbase/replication/regionserver/ReplicationSource.java
 (original)
+++ 
hbase/trunk/src/main/java/org/apache/hadoop/hbase/replication/regionserver/ReplicationSource.java
 Thu Sep 15 21:48:59 2011
@@ -240,7 +240,7 @@ public class ReplicationSource extends T
   public void run() {
     connectToPeers();
     // We were stopped while looping to connect to sinks, just abort
-    if (this.stopper.isStopped()) {
+    if (!this.isActive()) {
       return;
     }
     // delay this until we are in an asynchronous thread
@@ -265,7 +265,7 @@ public class ReplicationSource extends T
     }
     int sleepMultiplier = 1;
     // Loop until we close down
-    while (!stopper.isStopped() && this.running) {
+    while (isActive()) {
       // Sleep until replication is enabled again
       if (!this.replicating.get() || !this.sourceEnabled.get()) {
         if (sleepForRetries("Replication is disabled", sleepMultiplier)) {
@@ -348,7 +348,7 @@ public class ReplicationSource extends T
       // If we didn't get anything to replicate, or if we hit a IOE,
       // wait a bit and retry.
       // But if we need to stop, don't bother sleeping
-      if (!stopper.isStopped() && (gotIOE || currentNbEntries == 0)) {
+      if (this.isActive() && (gotIOE || currentNbEntries == 0)) {
         this.manager.logPositionAndCleanOldLogs(this.currentPath,
             this.peerClusterZnode, this.position, queueRecovered);
         if (sleepForRetries("Nothing to replicate", sleepMultiplier)) {
@@ -428,7 +428,8 @@ public class ReplicationSource extends T
 
   private void connectToPeers() {
     // Connect to peer cluster first, unless we have to stop
-    while (!this.stopper.isStopped() && this.currentPeers.size() == 0) {
+    while (this.isActive() && this.currentPeers.size() == 0) {
+
       try {
         chooseSinks();
         Thread.sleep(this.sleepForRetries);
@@ -586,7 +587,7 @@ public class ReplicationSource extends T
       LOG.warn("Was given 0 edits to ship");
       return;
     }
-    while (!this.stopper.isStopped()) {
+    while (this.isActive()) {
       try {
         HRegionInterface rrs = getRS();
         LOG.debug("Replicating " + currentNbEntries);
@@ -613,6 +614,7 @@ public class ReplicationSource extends T
         }
         try {
           boolean down;
+          // Spin while the slave is down and we're not asked to shutdown/close
           do {
             down = isSlaveDown();
             if (down) {
@@ -622,7 +624,7 @@ public class ReplicationSource extends T
                 chooseSinks();
               }
             }
-          } while (!this.stopper.isStopped() && down);
+          } while (this.isActive() && down );
         } catch (InterruptedException e) {
           LOG.debug("Interrupted while trying to contact the peer cluster");
         } catch (KeeperException e) {
@@ -742,6 +744,10 @@ public class ReplicationSource extends T
     this.sourceEnabled.set(status);
   }
 
+  private boolean isActive() {
+    return !this.stopper.isStopped() && this.running;
+  }
+
   /**
    * Comparator used to compare logs together based on their start time
    */


Reply via email to