Apache9 commented on code in PR #5247:
URL: https://github.com/apache/hbase/pull/5247#discussion_r1202502990


##########
hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/MemStoreFlusher.java:
##########
@@ -495,41 +505,54 @@ public int getFlushQueueSize() {
    * Only interrupt once it's done with a run through the work loop.
    */
   void interruptIfNecessary() {
-    lock.writeLock().lock();
+    lock.readLock().lock();

Review Comment:
   Mind explaining a bit why readLock is enough here?



##########
hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/MemStoreFlusher.java:
##########
@@ -495,41 +505,54 @@ public int getFlushQueueSize() {
    * Only interrupt once it's done with a run through the work loop.
    */
   void interruptIfNecessary() {
-    lock.writeLock().lock();
+    lock.readLock().lock();
     try {
-      for (FlushHandler flushHander : flushHandlers) {
-        if (flushHander != null) flushHander.interrupt();
+      for (FlushHandler flushHandler : flushHandlers) {
+        if (flushHandler != null) {
+          flushHandler.interrupt();
+        }
       }
     } finally {
-      lock.writeLock().unlock();
+      lock.readLock().unlock();
     }
   }
 
   synchronized void start(UncaughtExceptionHandler eh) {
-    ThreadFactory flusherThreadFactory = new ThreadFactoryBuilder()
+    this.flusherThreadFactory = new ThreadFactoryBuilder()
       .setNameFormat(server.getServerName().toShortString() + 
"-MemStoreFlusher-pool-%d")
       .setDaemon(true).setUncaughtExceptionHandler(eh).build();
-    for (int i = 0; i < flushHandlers.length; i++) {
-      flushHandlers[i] = new FlushHandler("MemStoreFlusher." + i);
-      flusherThreadFactory.newThread(flushHandlers[i]);
-      flushHandlers[i].start();
+    lock.readLock().lock();
+    try {
+      startFlushHandlerThreads(flushHandlers, 0, flushHandlers.length);
+    } finally {
+      lock.readLock().unlock();
     }
   }
 
   boolean isAlive() {
-    for (FlushHandler flushHander : flushHandlers) {
-      if (flushHander != null && flushHander.isAlive()) {
-        return true;
+    lock.readLock().lock();
+    try {
+      for (FlushHandler flushHandler : flushHandlers) {
+        if (flushHandler != null && flushHandler.isAlive()) {
+          return true;
+        }
       }
+      return false;
+    } finally {
+      lock.readLock().unlock();
     }
-    return false;
   }
 
   void join() {

Review Comment:
   The name is join but we call shutdown? Is this the expected behavior? If so 
I think we should change the method name?



##########
hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/MemStoreFlusher.java:
##########
@@ -924,4 +947,62 @@ public boolean equals(Object obj) {
       return compareTo(other) == 0;
     }
   }
+
+  private int getHandlerCount(Configuration conf) {
+    int handlerCount = conf.getInt("hbase.hstore.flusher.count", 2);
+    if (handlerCount < 1) {
+      LOG.warn(
+        "hbase.hstore.flusher.count was configed to {} which is less than 1, " 
+ "corrected to 1",
+        handlerCount);
+      handlerCount = 1;
+    }
+    return handlerCount;
+  }
+
+  @Override
+  public void onConfigurationChange(Configuration newConf) {
+    int newHandlerCount = getHandlerCount(newConf);
+    if (newHandlerCount != flushHandlers.length) {
+      LOG.info("update hbase.hstore.flusher.count from {} to {}", 
flushHandlers.length,
+        newHandlerCount);
+      lock.writeLock().lock();
+      try {
+        FlushHandler[] newFlushHandlers = new FlushHandler[newHandlerCount];
+        if (newHandlerCount > flushHandlers.length) {
+          System.arraycopy(flushHandlers, 0, newFlushHandlers, 0, 
flushHandlers.length);
+          startFlushHandlerThreads(newFlushHandlers, flushHandlers.length, 
newFlushHandlers.length);
+        } else {
+          System.arraycopy(flushHandlers, 0, newFlushHandlers, 0, 
newFlushHandlers.length);
+          stopFlushHandlerThreads(flushHandlers, newHandlerCount, 
flushHandlers.length);
+        }
+        flusherIdGen.compareAndSet(flushHandlers.length, 
newFlushHandlers.length);
+        this.flushHandlers = newFlushHandlers;
+      } finally {
+        lock.writeLock().unlock();
+      }
+    }
+  }
+
+  private void startFlushHandlerThreads(FlushHandler[] flushHandlers, int 
start, int end) {
+    if (flusherThreadFactory != null) {
+      for (int i = start; i < end; i++) {
+        flushHandlers[i] = new FlushHandler("MemStoreFlusher." + 
flusherIdGen.getAndIncrement());
+        flusherThreadFactory.newThread(flushHandlers[i]);
+        flushHandlers[i].start();
+      }
+    }
+  }
+
+  private void stopFlushHandlerThreads(FlushHandler[] flushHandlers, int 
start, int end) {
+    for (int i = start; i < end; i++) {
+      flushHandlers[i].shutdown();

Review Comment:
   So the shutdown here will not interrupt the current flush operation if any 
right?



##########
hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/MemStoreFlusher.java:
##########
@@ -924,4 +947,62 @@ public boolean equals(Object obj) {
       return compareTo(other) == 0;
     }
   }
+
+  private int getHandlerCount(Configuration conf) {
+    int handlerCount = conf.getInt("hbase.hstore.flusher.count", 2);
+    if (handlerCount < 1) {
+      LOG.warn(
+        "hbase.hstore.flusher.count was configed to {} which is less than 1, " 
+ "corrected to 1",
+        handlerCount);
+      handlerCount = 1;
+    }
+    return handlerCount;
+  }
+
+  @Override
+  public void onConfigurationChange(Configuration newConf) {
+    int newHandlerCount = getHandlerCount(newConf);
+    if (newHandlerCount != flushHandlers.length) {
+      LOG.info("update hbase.hstore.flusher.count from {} to {}", 
flushHandlers.length,
+        newHandlerCount);
+      lock.writeLock().lock();
+      try {
+        FlushHandler[] newFlushHandlers = new FlushHandler[newHandlerCount];
+        if (newHandlerCount > flushHandlers.length) {
+          System.arraycopy(flushHandlers, 0, newFlushHandlers, 0, 
flushHandlers.length);

Review Comment:
   Use Arrays.copyOf, it can handle both cases, where the new length is greater 
or less.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to