This is an automated email from the ASF dual-hosted git repository. ericpai pushed a commit to branch bugfix/iotdb-3958 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 53e4e7d50560a4e6ce72e96ba49fdeabd003d0c6 Author: ericpai <[email protected]> AuthorDate: Tue Jul 26 14:21:29 2022 +0800 [IOTDB-3958] Add error logging in thread pool --- client-go | 2 +- .../{WrappedRunnable.java => WrappedCallable.java} | 32 ++++++++++++---------- .../iotdb/commons/concurrent/WrappedRunnable.java | 24 ++++++++-------- .../threadpool/ScheduledExecutorUtil.java | 7 +++++ .../WrappedScheduledExecutorService.java | 31 +++++++++++++-------- .../WrappedSingleThreadExecutorService.java | 23 ++++++++++------ .../WrappedSingleThreadScheduledExecutor.java | 31 +++++++++++++-------- .../threadpool/WrappedThreadPoolExecutor.java | 25 +++++++++++++++++ .../apache/iotdb/db/engine/flush/FlushManager.java | 2 +- .../apache/iotdb/db/engine/settle/SettleTask.java | 2 +- .../iotdb/db/engine/upgrade/UpgradeTask.java | 2 +- .../db/query/dataset/NonAlignEngineDataSet.java | 2 +- .../dataset/RawQueryDataSetWithoutValueFilter.java | 2 +- 13 files changed, 122 insertions(+), 63 deletions(-) diff --git a/client-go b/client-go index 84b8d45829..3208ea6a77 160000 --- a/client-go +++ b/client-go @@ -1 +1 @@ -Subproject commit 84b8d45829d846440a3246400e7bc5e39587dcb5 +Subproject commit 3208ea6a77b980b455d96da47240951fd65b92e7 diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedRunnable.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedCallable.java similarity index 56% copy from node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedRunnable.java copy to node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedCallable.java index 738c31b785..ae42a0dc4f 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedRunnable.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedCallable.java @@ -18,29 +18,33 @@ */ package org.apache.iotdb.commons.concurrent; -import com.google.common.base.Throwables; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; -public abstract class WrappedRunnable implements Runnable { +import java.util.concurrent.Callable; - private static final Logger LOGGER = LoggerFactory.getLogger(WrappedRunnable.class); +/** A wrapper for {@link Callable} logging errors when uncaught exception is thrown. */ +public abstract class WrappedCallable<V> implements Callable<V> { @Override - public final void run() { + public final V call() { try { - runMayThrow(); + return callMayThrow(); } catch (Exception e) { - LOGGER.error(e.getMessage(), e); - throw propagate(e); + throw ScheduledExecutorUtil.propagate(e); } } - public abstract void runMayThrow() throws Exception; + public abstract V callMayThrow() throws Exception; - @SuppressWarnings("squid:S112") - private static RuntimeException propagate(Throwable throwable) { - Throwables.throwIfUnchecked(throwable); - throw new RuntimeException(throwable); + public static <V> Callable<V> wrap(Callable<V> callable) { + if (callable instanceof WrappedCallable) { + return callable; + } + return new WrappedCallable<V>() { + @Override + public V callMayThrow() throws Exception { + return callable.call(); + } + }; } } diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedRunnable.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedRunnable.java index 738c31b785..a91278639d 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedRunnable.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/WrappedRunnable.java @@ -18,29 +18,31 @@ */ package org.apache.iotdb.commons.concurrent; -import com.google.common.base.Throwables; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; +/** A wrapper for {@link Runnable} logging errors when uncaught exception is thrown. */ public abstract class WrappedRunnable implements Runnable { - private static final Logger LOGGER = LoggerFactory.getLogger(WrappedRunnable.class); - @Override public final void run() { try { runMayThrow(); } catch (Exception e) { - LOGGER.error(e.getMessage(), e); - throw propagate(e); + throw ScheduledExecutorUtil.propagate(e); } } public abstract void runMayThrow() throws Exception; - @SuppressWarnings("squid:S112") - private static RuntimeException propagate(Throwable throwable) { - Throwables.throwIfUnchecked(throwable); - throw new RuntimeException(throwable); + public static Runnable wrap(Runnable runnable) { + if (runnable instanceof WrappedRunnable) { + return runnable; + } + return new WrappedRunnable() { + @Override + public void runMayThrow() { + runnable.run(); + } + }; } } diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/ScheduledExecutorUtil.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/ScheduledExecutorUtil.java index 1ed2993a12..36b4f189ad 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/ScheduledExecutorUtil.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/ScheduledExecutorUtil.java @@ -18,6 +18,7 @@ */ package org.apache.iotdb.commons.concurrent.threadpool; +import com.google.common.base.Throwables; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -185,4 +186,10 @@ public class ScheduledExecutorUtil { delay, unit); } + + public static RuntimeException propagate(Throwable throwable) { + logger.error("Run thread failed", throwable); + Throwables.throwIfUnchecked(throwable); + throw new RuntimeException(throwable); + } } diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedScheduledExecutorService.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedScheduledExecutorService.java index a362983785..1dbd7f369c 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedScheduledExecutorService.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedScheduledExecutorService.java @@ -19,6 +19,8 @@ package org.apache.iotdb.commons.concurrent.threadpool; +import org.apache.iotdb.commons.concurrent.WrappedCallable; +import org.apache.iotdb.commons.concurrent.WrappedRunnable; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.service.JMXService; @@ -33,6 +35,7 @@ import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.stream.Collectors; public class WrappedScheduledExecutorService implements ScheduledExecutorService, WrappedScheduledExecutorServiceMBean { @@ -49,26 +52,26 @@ public class WrappedScheduledExecutorService @Override public ScheduledFuture<?> schedule(Runnable command, long delay, TimeUnit unit) { - return service.schedule(command, delay, unit); + return service.schedule(WrappedRunnable.wrap(command), delay, unit); } @Override public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay, TimeUnit unit) { - return service.schedule(callable, delay, unit); + return service.schedule(WrappedCallable.wrap(callable), delay, unit); } @Override @SuppressWarnings("unsafeThreadSchedule") public ScheduledFuture<?> scheduleAtFixedRate( Runnable command, long initialDelay, long period, TimeUnit unit) { - return service.scheduleAtFixedRate(command, initialDelay, period, unit); + return service.scheduleAtFixedRate(WrappedRunnable.wrap(command), initialDelay, period, unit); } @Override @SuppressWarnings("unsafeThreadSchedule") public ScheduledFuture<?> scheduleWithFixedDelay( Runnable command, long initialDelay, long delay, TimeUnit unit) { - return service.scheduleWithFixedDelay(command, initialDelay, delay, unit); + return service.scheduleWithFixedDelay(WrappedRunnable.wrap(command), initialDelay, delay, unit); } @Override @@ -100,47 +103,51 @@ public class WrappedScheduledExecutorService @Override public <T> Future<T> submit(Callable<T> task) { - return service.submit(task); + return service.submit(WrappedCallable.wrap(task)); } @Override public <T> Future<T> submit(Runnable task, T result) { - return service.submit(task, result); + return service.submit(WrappedRunnable.wrap(task), result); } @Override public Future<?> submit(Runnable task) { - return service.submit(task); + return service.submit(WrappedRunnable.wrap(task)); } @Override public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException { - return service.invokeAll(tasks); + return service.invokeAll( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList())); } @Override public <T> List<Future<T>> invokeAll( Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException { - return service.invokeAll(tasks, timeout, unit); + return service.invokeAll( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList()), timeout, unit); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException { - return service.invokeAny(tasks); + return service.invokeAny( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList())); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { - return service.invokeAny(tasks, timeout, unit); + return service.invokeAny( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList()), timeout, unit); } @Override public void execute(Runnable command) { - service.execute(command); + service.execute(WrappedRunnable.wrap(command)); } @Override diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadExecutorService.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadExecutorService.java index 9035c4ba35..f597669815 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadExecutorService.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadExecutorService.java @@ -19,6 +19,8 @@ package org.apache.iotdb.commons.concurrent.threadpool; +import org.apache.iotdb.commons.concurrent.WrappedCallable; +import org.apache.iotdb.commons.concurrent.WrappedRunnable; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.service.JMXService; @@ -30,6 +32,7 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.stream.Collectors; public class WrappedSingleThreadExecutorService implements ExecutorService, WrappedSingleThreadExecutorServiceMBean { @@ -74,46 +77,50 @@ public class WrappedSingleThreadExecutorService @Override public <T> Future<T> submit(Callable<T> task) { - return service.submit(task); + return service.submit(WrappedCallable.wrap(task)); } @Override public <T> Future<T> submit(Runnable task, T result) { - return service.submit(task, result); + return service.submit(WrappedRunnable.wrap(task), result); } @Override public Future<?> submit(Runnable task) { - return service.submit(task); + return service.submit(WrappedRunnable.wrap(task)); } @Override public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException { - return service.invokeAll(tasks); + return service.invokeAll( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList())); } @Override public <T> List<Future<T>> invokeAll( Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException { - return service.invokeAll(tasks, timeout, unit); + return service.invokeAll( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList()), timeout, unit); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException { - return service.invokeAny(tasks); + return service.invokeAny( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList())); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { - return service.invokeAny(tasks, timeout, unit); + return service.invokeAny( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList()), timeout, unit); } @Override public void execute(Runnable command) { - service.execute(command); + service.execute(WrappedRunnable.wrap(command)); } } diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadScheduledExecutor.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadScheduledExecutor.java index 40d5e7183c..b8238a0405 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadScheduledExecutor.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedSingleThreadScheduledExecutor.java @@ -19,6 +19,8 @@ package org.apache.iotdb.commons.concurrent.threadpool; +import org.apache.iotdb.commons.concurrent.WrappedCallable; +import org.apache.iotdb.commons.concurrent.WrappedRunnable; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.service.JMXService; @@ -31,6 +33,7 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.stream.Collectors; public class WrappedSingleThreadScheduledExecutor implements ScheduledExecutorService, WrappedSingleThreadScheduledExecutorMBean { @@ -47,26 +50,26 @@ public class WrappedSingleThreadScheduledExecutor @Override public ScheduledFuture<?> schedule(Runnable command, long delay, TimeUnit unit) { - return service.schedule(command, delay, unit); + return service.schedule(WrappedRunnable.wrap(command), delay, unit); } @Override public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay, TimeUnit unit) { - return service.schedule(callable, delay, unit); + return service.schedule(WrappedCallable.wrap(callable), delay, unit); } @Override @SuppressWarnings("unsafeThreadSchedule") public ScheduledFuture<?> scheduleAtFixedRate( Runnable command, long initialDelay, long period, TimeUnit unit) { - return service.scheduleAtFixedRate(command, initialDelay, period, unit); + return service.scheduleAtFixedRate(WrappedRunnable.wrap(command), initialDelay, period, unit); } @Override @SuppressWarnings("unsafeThreadSchedule") public ScheduledFuture<?> scheduleWithFixedDelay( Runnable command, long initialDelay, long delay, TimeUnit unit) { - return service.scheduleWithFixedDelay(command, initialDelay, delay, unit); + return service.scheduleWithFixedDelay(WrappedRunnable.wrap(command), initialDelay, delay, unit); } @Override @@ -98,46 +101,50 @@ public class WrappedSingleThreadScheduledExecutor @Override public <T> Future<T> submit(Callable<T> task) { - return service.submit(task); + return service.submit(WrappedCallable.wrap(task)); } @Override public <T> Future<T> submit(Runnable task, T result) { - return service.submit(task, result); + return service.submit(WrappedRunnable.wrap(task), result); } @Override public Future<?> submit(Runnable task) { - return service.submit(task); + return service.submit(WrappedRunnable.wrap(task)); } @Override public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException { - return service.invokeAll(tasks); + return service.invokeAll( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList())); } @Override public <T> List<Future<T>> invokeAll( Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException { - return service.invokeAll(tasks, timeout, unit); + return service.invokeAll( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList()), timeout, unit); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException { - return service.invokeAny(tasks); + return service.invokeAny( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList())); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { - return service.invokeAny(tasks, timeout, unit); + return service.invokeAny( + tasks.stream().map(WrappedCallable::wrap).collect(Collectors.toList()), timeout, unit); } @Override public void execute(Runnable command) { - service.execute(command); + service.execute(WrappedRunnable.wrap(command)); } } diff --git a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedThreadPoolExecutor.java b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedThreadPoolExecutor.java index 834f5d5500..4538603e60 100644 --- a/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedThreadPoolExecutor.java +++ b/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/threadpool/WrappedThreadPoolExecutor.java @@ -23,14 +23,21 @@ import org.apache.iotdb.commons.concurrent.IoTThreadFactory; import org.apache.iotdb.commons.conf.IoTDBConstant; import org.apache.iotdb.commons.service.JMXService; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import java.util.List; import java.util.concurrent.BlockingQueue; +import java.util.concurrent.CancellationException; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Future; import java.util.concurrent.ThreadFactory; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; public class WrappedThreadPoolExecutor extends ThreadPoolExecutor implements WrappedThreadPoolExecutorMBean { + private static final Logger logger = LoggerFactory.getLogger(WrappedThreadPoolExecutor.class); private final String mbeanName; public WrappedThreadPoolExecutor( @@ -79,4 +86,22 @@ public class WrappedThreadPoolExecutor extends ThreadPoolExecutor public int getQueueLength() { return getQueue().size(); } + + protected void afterExecute(Runnable r, Throwable t) { + super.afterExecute(r, t); + if (t == null && r instanceof Future<?>) { + try { + ((Future<?>) r).get(); + } catch (CancellationException ce) { + t = ce; + } catch (ExecutionException ee) { + t = ee.getCause(); + } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); + } + } + if (t != null) { + logger.error("Exception in thread pool {}", mbeanName, t); + } + } } diff --git a/server/src/main/java/org/apache/iotdb/db/engine/flush/FlushManager.java b/server/src/main/java/org/apache/iotdb/db/engine/flush/FlushManager.java index d0455e0a7a..d9374c863c 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/flush/FlushManager.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/flush/FlushManager.java @@ -18,7 +18,7 @@ */ package org.apache.iotdb.db.engine.flush; -import org.apache.iotdb.commons.concurrent.WrappedRunnable; +import org.apache.iotdb.commons.concurrent.threadpool.WrappedRunnable; import org.apache.iotdb.commons.exception.StartupException; import org.apache.iotdb.commons.service.IService; import org.apache.iotdb.commons.service.JMXService; diff --git a/server/src/main/java/org/apache/iotdb/db/engine/settle/SettleTask.java b/server/src/main/java/org/apache/iotdb/db/engine/settle/SettleTask.java index 50e431c4de..cac5f1e147 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/settle/SettleTask.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/settle/SettleTask.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.engine.settle; -import org.apache.iotdb.commons.concurrent.WrappedRunnable; +import org.apache.iotdb.commons.concurrent.threadpool.WrappedRunnable; import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.db.engine.settle.SettleLog.SettleCheckStatus; import org.apache.iotdb.db.engine.storagegroup.TsFileResource; diff --git a/server/src/main/java/org/apache/iotdb/db/engine/upgrade/UpgradeTask.java b/server/src/main/java/org/apache/iotdb/db/engine/upgrade/UpgradeTask.java index 4ec2b14aaa..7526011f8f 100644 --- a/server/src/main/java/org/apache/iotdb/db/engine/upgrade/UpgradeTask.java +++ b/server/src/main/java/org/apache/iotdb/db/engine/upgrade/UpgradeTask.java @@ -18,7 +18,7 @@ */ package org.apache.iotdb.db.engine.upgrade; -import org.apache.iotdb.commons.concurrent.WrappedRunnable; +import org.apache.iotdb.commons.concurrent.threadpool.WrappedRunnable; import org.apache.iotdb.db.conf.directories.DirectoryManager; import org.apache.iotdb.db.engine.storagegroup.TsFileResource; import org.apache.iotdb.db.service.UpgradeSevice; diff --git a/server/src/main/java/org/apache/iotdb/db/query/dataset/NonAlignEngineDataSet.java b/server/src/main/java/org/apache/iotdb/db/query/dataset/NonAlignEngineDataSet.java index 721ddc2df4..fc73e990af 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/dataset/NonAlignEngineDataSet.java +++ b/server/src/main/java/org/apache/iotdb/db/query/dataset/NonAlignEngineDataSet.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.query.dataset; -import org.apache.iotdb.commons.concurrent.WrappedRunnable; +import org.apache.iotdb.commons.concurrent.threadpool.WrappedRunnable; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.db.query.control.QueryTimeManager; import org.apache.iotdb.db.query.pool.RawQueryReadTaskPoolManager; diff --git a/server/src/main/java/org/apache/iotdb/db/query/dataset/RawQueryDataSetWithoutValueFilter.java b/server/src/main/java/org/apache/iotdb/db/query/dataset/RawQueryDataSetWithoutValueFilter.java index 2184d74697..4d5356292d 100644 --- a/server/src/main/java/org/apache/iotdb/db/query/dataset/RawQueryDataSetWithoutValueFilter.java +++ b/server/src/main/java/org/apache/iotdb/db/query/dataset/RawQueryDataSetWithoutValueFilter.java @@ -19,7 +19,7 @@ package org.apache.iotdb.db.query.dataset; -import org.apache.iotdb.commons.concurrent.WrappedRunnable; +import org.apache.iotdb.commons.concurrent.threadpool.WrappedRunnable; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.db.conf.IoTDBDescriptor; import org.apache.iotdb.db.metadata.path.AlignedPath;
