This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new ec1b01adf0 [INLONG-9591][Manager] Support printing thread status
before submitting tasks (#9593)
ec1b01adf0 is described below
commit ec1b01adf057bfb5cb850da9782ea9c08c6f2a9a
Author: fuweng11 <[email protected]>
AuthorDate: Tue Jan 23 10:20:40 2024 +0800
[INLONG-9591][Manager] Support printing thread status before submitting
tasks (#9593)
---
.../threadPool/VisiableThreadPoolTaskExecutor.java | 72 ++++++++++++++++++++++
.../service/group/InlongGroupProcessService.java | 4 +-
.../service/stream/InlongStreamProcessService.java | 4 +-
3 files changed, 76 insertions(+), 4 deletions(-)
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/threadPool/VisiableThreadPoolTaskExecutor.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/threadPool/VisiableThreadPoolTaskExecutor.java
new file mode 100644
index 0000000000..e9803e74e8
--- /dev/null
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/threadPool/VisiableThreadPoolTaskExecutor.java
@@ -0,0 +1,72 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.manager.common.threadPool;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.Callable;
+import java.util.concurrent.Future;
+import java.util.concurrent.RejectedExecutionHandler;
+import java.util.concurrent.ThreadFactory;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+public class VisiableThreadPoolTaskExecutor extends ThreadPoolExecutor {
+
+ private static final Logger logger =
+ LoggerFactory.getLogger(VisiableThreadPoolTaskExecutor.class);
+
+ public VisiableThreadPoolTaskExecutor(int corePoolSize, int
maximumPoolSize, long keepAliveTime, TimeUnit unit,
+ BlockingQueue<Runnable> workQueue, ThreadFactory threadFactory,
+ RejectedExecutionHandler handler) {
+ super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue,
threadFactory, handler);
+ }
+
+ private void showThreadPoolInfo(String prefix) {
+ logger.info(
+ "Current thread pool class = {}, opType = {}, taskCount = {},
completedTaskCount = {}, activeCount = {}, poolSize = {}, queueSize = {}",
+ this.getThreadFactory().getClass(),
+ prefix,
+ this.getTaskCount(),
+ this.getCompletedTaskCount(),
+ this.getActiveCount(),
+ this.getPoolSize(),
+ this.getQueue().size());
+ }
+
+ @Override
+ public void execute(Runnable task) {
+ showThreadPoolInfo("execute");
+ super.execute(task);
+ }
+
+ @Override
+ public Future<?> submit(Runnable task) {
+ showThreadPoolInfo("submit");
+ return super.submit(task);
+ }
+
+ @Override
+ public <T> Future<T> submit(Callable<T> task) {
+ showThreadPoolInfo("submit");
+ return super.submit(task);
+ }
+
+}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupProcessService.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupProcessService.java
index b974e5f7ca..a1460e4ed4 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupProcessService.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupProcessService.java
@@ -26,6 +26,7 @@ import org.apache.inlong.manager.common.enums.TaskStatus;
import org.apache.inlong.manager.common.enums.TenantUserTypeEnum;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
+import
org.apache.inlong.manager.common.threadPool.VisiableThreadPoolTaskExecutor;
import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.dao.entity.InlongGroupEntity;
import org.apache.inlong.manager.dao.entity.WorkflowProcessEntity;
@@ -56,7 +57,6 @@ import java.util.Comparator;
import java.util.List;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
-import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy;
import java.util.concurrent.TimeUnit;
@@ -73,7 +73,7 @@ public class InlongGroupProcessService {
private static final Logger LOGGER =
LoggerFactory.getLogger(InlongGroupProcessService.class);
- private static final ExecutorService EXECUTOR_SERVICE = new
ThreadPoolExecutor(
+ private static final ExecutorService EXECUTOR_SERVICE = new
VisiableThreadPoolTaskExecutor(
CORE_POOL_SIZE,
MAX_POOL_SIZE,
ALIVE_TIME_MS,
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamProcessService.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamProcessService.java
index 367e82d718..b37fd7b5b6 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamProcessService.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/stream/InlongStreamProcessService.java
@@ -24,6 +24,7 @@ import org.apache.inlong.manager.common.enums.ProcessName;
import org.apache.inlong.manager.common.enums.ProcessStatus;
import org.apache.inlong.manager.common.enums.StreamStatus;
import org.apache.inlong.manager.common.exceptions.BusinessException;
+import
org.apache.inlong.manager.common.threadPool.VisiableThreadPoolTaskExecutor;
import org.apache.inlong.manager.common.util.Preconditions;
import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamInfo;
@@ -41,7 +42,6 @@ import org.springframework.stereotype.Service;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
-import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy;
import java.util.concurrent.TimeUnit;
@@ -57,7 +57,7 @@ import static
org.apache.inlong.manager.common.consts.InlongConstants.QUEUE_SIZE
@Service
public class InlongStreamProcessService {
- private static final ExecutorService EXECUTOR_SERVICE = new
ThreadPoolExecutor(
+ private static final ExecutorService EXECUTOR_SERVICE = new
VisiableThreadPoolTaskExecutor(
CORE_POOL_SIZE,
MAX_POOL_SIZE,
ALIVE_TIME_MS,