This is an automated email from the ASF dual-hosted git repository.

rong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 9d9d6ea6eaa Pipe: add seperated thread pool for phantom reference 
clean job (#13813)
9d9d6ea6eaa is described below

commit 9d9d6ea6eaa937422cbaeb2ff05be6df697ab391
Author: V_Galaxy <[email protected]>
AuthorDate: Wed Oct 23 19:49:35 2024 +0800

    Pipe: add seperated thread pool for phantom reference clean job (#13813)
---
 .../agent/runtime/PipeConfigNodeRuntimeAgent.java  |  13 +++
 .../ref/PipeConfigNodePhantomReferenceManager.java |   2 +-
 .../agent/runtime/PipeDataNodeRuntimeAgent.java    |  13 +++
 .../ref/PipeDataNodePhantomReferenceManager.java   |   2 +-
 .../resource/tsfile/PipeTsFileResourceManager.java |   4 +-
 .../iotdb/commons/concurrent/ThreadName.java       |   2 +
 .../apache/iotdb/commons/conf/CommonConfig.java    |   2 +-
 ...java => AbstractPipePeriodicalJobExecutor.java} |  38 ++++----
 .../agent/runtime/PipePeriodicalJobExecutor.java   | 101 ++-------------------
 .../PipePeriodicalPhantomReferenceCleaner.java}    |  24 ++---
 .../resource/ref/PipePhantomReferenceManager.java  |  20 +++-
 11 files changed, 87 insertions(+), 134 deletions(-)

diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/runtime/PipeConfigNodeRuntimeAgent.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/runtime/PipeConfigNodeRuntimeAgent.java
index 99dbbf1f56d..b89e401d5c7 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/runtime/PipeConfigNodeRuntimeAgent.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/runtime/PipeConfigNodeRuntimeAgent.java
@@ -23,6 +23,7 @@ import 
org.apache.iotdb.commons.exception.IllegalPathException;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeCriticalException;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeException;
 import org.apache.iotdb.commons.pipe.agent.runtime.PipePeriodicalJobExecutor;
+import 
org.apache.iotdb.commons.pipe.agent.runtime.PipePeriodicalPhantomReferenceCleaner;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
 import org.apache.iotdb.commons.pipe.config.PipeConfig;
 import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
@@ -49,6 +50,9 @@ public class PipeConfigNodeRuntimeAgent implements IService {
   private final PipePeriodicalJobExecutor pipePeriodicalJobExecutor =
       new PipePeriodicalJobExecutor();
 
+  private final PipePeriodicalPhantomReferenceCleaner 
pipePeriodicalPhantomReferenceCleaner =
+      new PipePeriodicalPhantomReferenceCleaner();
+
   @Override
   public synchronized void start() {
     PipeConfig.getInstance().printAllConfigs();
@@ -65,6 +69,10 @@ public class PipeConfigNodeRuntimeAgent implements IService {
     // Start periodical job executor
     pipePeriodicalJobExecutor.start();
 
+    if (PipeConfig.getInstance().getPipeEventReferenceTrackingEnabled()) {
+      pipePeriodicalPhantomReferenceCleaner.start();
+    }
+
     isShutdown.set(false);
     LOGGER.info("PipeRuntimeConfigNodeAgent started");
   }
@@ -159,4 +167,9 @@ public class PipeConfigNodeRuntimeAgent implements IService 
{
   public void registerPeriodicalJob(String id, Runnable periodicalJob, long 
intervalInSeconds) {
     pipePeriodicalJobExecutor.register(id, periodicalJob, intervalInSeconds);
   }
+
+  public void registerPhantomReferenceCleanJob(
+      String id, Runnable periodicalJob, long intervalInSeconds) {
+    pipePeriodicalPhantomReferenceCleaner.register(id, periodicalJob, 
intervalInSeconds);
+  }
 }
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/resource/ref/PipeConfigNodePhantomReferenceManager.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/resource/ref/PipeConfigNodePhantomReferenceManager.java
index 13298564ce4..b867163e47e 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/resource/ref/PipeConfigNodePhantomReferenceManager.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/resource/ref/PipeConfigNodePhantomReferenceManager.java
@@ -29,7 +29,7 @@ public class PipeConfigNodePhantomReferenceManager extends 
PipePhantomReferenceM
     super();
 
     PipeConfigNodeAgent.runtime()
-        .registerPeriodicalJob(
+        .registerPhantomReferenceCleanJob(
             "PipePhantomReferenceManager#gcHook()",
             // NOTE: lambda CAN NOT be replaced with method reference
             () -> super.gcHook(),
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeDataNodeRuntimeAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeDataNodeRuntimeAgent.java
index 7b2a6d69dcc..057000dd965 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeDataNodeRuntimeAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeDataNodeRuntimeAgent.java
@@ -26,6 +26,7 @@ import org.apache.iotdb.commons.exception.StartupException;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeCriticalException;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeException;
 import org.apache.iotdb.commons.pipe.agent.runtime.PipePeriodicalJobExecutor;
+import 
org.apache.iotdb.commons.pipe.agent.runtime.PipePeriodicalPhantomReferenceCleaner;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
 import org.apache.iotdb.commons.pipe.config.PipeConfig;
 import org.apache.iotdb.commons.pipe.event.EnrichedEvent;
@@ -62,6 +63,9 @@ public class PipeDataNodeRuntimeAgent implements IService {
   private final PipePeriodicalJobExecutor pipePeriodicalJobExecutor =
       new PipePeriodicalJobExecutor();
 
+  private final PipePeriodicalPhantomReferenceCleaner 
pipePeriodicalPhantomReferenceCleaner =
+      new PipePeriodicalPhantomReferenceCleaner();
+
   //////////////////////////// System Service Interface 
////////////////////////////
 
   public synchronized void preparePipeResources(
@@ -87,6 +91,10 @@ public class PipeDataNodeRuntimeAgent implements IService {
         PipeConfig.getInstance().getPipeStuckRestartIntervalSeconds());
     pipePeriodicalJobExecutor.start();
 
+    if (PipeConfig.getInstance().getPipeEventReferenceTrackingEnabled()) {
+      pipePeriodicalPhantomReferenceCleaner.start();
+    }
+
     isShutdown.set(false);
   }
 
@@ -225,4 +233,9 @@ public class PipeDataNodeRuntimeAgent implements IService {
   public void clearPeriodicalJobExecutor() {
     pipePeriodicalJobExecutor.clear();
   }
+
+  public void registerPhantomReferenceCleanJob(
+      String id, Runnable periodicalJob, long intervalInSeconds) {
+    pipePeriodicalPhantomReferenceCleaner.register(id, periodicalJob, 
intervalInSeconds);
+  }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/ref/PipeDataNodePhantomReferenceManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/ref/PipeDataNodePhantomReferenceManager.java
index 06faed209a9..991cfebe8fa 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/ref/PipeDataNodePhantomReferenceManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/ref/PipeDataNodePhantomReferenceManager.java
@@ -29,7 +29,7 @@ public class PipeDataNodePhantomReferenceManager extends 
PipePhantomReferenceMan
     super();
 
     PipeDataNodeAgent.runtime()
-        .registerPeriodicalJob(
+        .registerPhantomReferenceCleanJob(
             "PipePhantomReferenceManager#gcHook()",
             // NOTE: lambda CAN NOT be replaced with method reference
             () -> super.gcHook(),
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResourceManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResourceManager.java
index 02ba8552ee2..bb8bc948de8 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResourceManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/tsfile/PipeTsFileResourceManager.java
@@ -20,7 +20,6 @@
 package org.apache.iotdb.db.pipe.resource.tsfile;
 
 import org.apache.iotdb.commons.conf.IoTDBConstant;
-import org.apache.iotdb.commons.pipe.agent.runtime.PipePeriodicalJobExecutor;
 import org.apache.iotdb.commons.pipe.config.PipeConfig;
 import org.apache.iotdb.commons.utils.FileUtils;
 import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
@@ -61,7 +60,8 @@ public class PipeTsFileResourceManager {
 
   private void tryTtlCheck() {
     try {
-      final long timeout = PipePeriodicalJobExecutor.getMinIntervalSeconds() 
>> 1;
+      final long timeout =
+          
PipeConfig.getInstance().getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds()
 >> 1;
       if (lock.tryLock(timeout, TimeUnit.SECONDS)) {
         try {
           ttlCheck();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
index 0c9ac1eeffa..cd3416f6baf 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/concurrent/ThreadName.java
@@ -140,6 +140,8 @@ public enum ThreadName {
   PIPE_RUNTIME_HEARTBEAT("Pipe-Runtime-Heartbeat"),
   PIPE_RUNTIME_PROCEDURE_SUBMITTER("Pipe-Runtime-Procedure-Submitter"),
   PIPE_RUNTIME_PERIODICAL_JOB_EXECUTOR("Pipe-Runtime-Periodical-Job-Executor"),
+  PIPE_RUNTIME_PERIODICAL_PHANTOM_REFERENCE_CLEANER(
+      "Pipe-Runtime-Periodical-Phantom-Reference-Cleaner"),
   PIPE_ASYNC_CONNECTOR_CLIENT_POOL("Pipe-Async-Connector-Client-Pool"),
   PIPE_RECEIVER_AIR_GAP_AGENT("Pipe-Receiver-Air-Gap-Agent"),
   SUBSCRIPTION_EXECUTOR_POOL("Subscription-Executor-Pool"),
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
index 4d29f22482a..b0f67d6493e 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/conf/CommonConfig.java
@@ -274,7 +274,7 @@ public class CommonConfig {
 
   private float subscriptionCacheMemoryUsagePercentage = 0.2F;
 
-  private boolean pipeEventReferenceTrackingEnabled = false; // TODO: enable 
later
+  private boolean pipeEventReferenceTrackingEnabled = true;
   private long pipeEventReferenceEliminateIntervalSeconds = 10;
 
   private int subscriptionSubtaskExecutorMaxThreadNum =
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalJobExecutor.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/AbstractPipePeriodicalJobExecutor.java
similarity index 75%
copy from 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalJobExecutor.java
copy to 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/AbstractPipePeriodicalJobExecutor.java
index ca972f1cdd0..c6e24b5f475 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalJobExecutor.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/AbstractPipePeriodicalJobExecutor.java
@@ -19,11 +19,8 @@
 
 package org.apache.iotdb.commons.pipe.agent.runtime;
 
-import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
-import org.apache.iotdb.commons.concurrent.ThreadName;
 import org.apache.iotdb.commons.concurrent.WrappedRunnable;
 import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil;
-import org.apache.iotdb.commons.pipe.config.PipeConfig;
 import org.apache.iotdb.commons.utils.TestOnly;
 
 import org.apache.tsfile.utils.Pair;
@@ -40,16 +37,13 @@ import java.util.concurrent.TimeUnit;
  * Single thread to execute pipe periodical jobs on DataNode or ConfigNode. 
This is for limiting the
  * thread num on the DataNode or ConfigNode instance.
  */
-public class PipePeriodicalJobExecutor {
+public abstract class AbstractPipePeriodicalJobExecutor {
 
-  private static final Logger LOGGER = 
LoggerFactory.getLogger(PipePeriodicalJobExecutor.class);
+  private static final Logger LOGGER =
+      LoggerFactory.getLogger(AbstractPipePeriodicalJobExecutor.class);
 
-  private static final ScheduledExecutorService PERIODICAL_JOB_EXECUTOR =
-      IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(
-          ThreadName.PIPE_RUNTIME_PERIODICAL_JOB_EXECUTOR.getName());
-
-  private static final long MIN_INTERVAL_SECONDS =
-      
PipeConfig.getInstance().getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds();
+  private final ScheduledExecutorService executorService;
+  private final long minIntervalSeconds;
 
   private long rounds;
   private Future<?> executorFuture;
@@ -57,6 +51,12 @@ public class PipePeriodicalJobExecutor {
   // <Periodical job, Interval in rounds>
   private final List<Pair<WrappedRunnable, Long>> periodicalJobs = new 
CopyOnWriteArrayList<>();
 
+  public AbstractPipePeriodicalJobExecutor(
+      final ScheduledExecutorService executorService, final long 
minIntervalSeconds) {
+    this.executorService = executorService;
+    this.minIntervalSeconds = minIntervalSeconds;
+  }
+
   public void register(String id, Runnable periodicalJob, long 
intervalInSeconds) {
     periodicalJobs.add(
         new Pair<>(
@@ -70,11 +70,11 @@ public class PipePeriodicalJobExecutor {
                 }
               }
             },
-            Math.max(intervalInSeconds / MIN_INTERVAL_SECONDS, 1)));
+            Math.max(intervalInSeconds / minIntervalSeconds, 1)));
     LOGGER.info(
         "Pipe periodical job {} is registered successfully. Interval: {} 
seconds.",
         id,
-        Math.max(intervalInSeconds / MIN_INTERVAL_SECONDS, 1) * 
MIN_INTERVAL_SECONDS);
+        Math.max(intervalInSeconds / minIntervalSeconds, 1) * 
minIntervalSeconds);
   }
 
   public synchronized void start() {
@@ -83,16 +83,16 @@ public class PipePeriodicalJobExecutor {
 
       executorFuture =
           ScheduledExecutorUtil.safelyScheduleWithFixedDelay(
-              PERIODICAL_JOB_EXECUTOR,
+              executorService,
               this::execute,
-              MIN_INTERVAL_SECONDS,
-              MIN_INTERVAL_SECONDS,
+              minIntervalSeconds,
+              minIntervalSeconds,
               TimeUnit.SECONDS);
       LOGGER.info("Pipe periodical job executor is started successfully.");
     }
   }
 
-  private void execute() {
+  protected void execute() {
     ++rounds;
 
     for (final Pair<WrappedRunnable, Long> periodicalJob : periodicalJobs) {
@@ -115,8 +115,4 @@ public class PipePeriodicalJobExecutor {
     periodicalJobs.clear();
     LOGGER.info("All pipe periodical jobs are cleared successfully.");
   }
-
-  public static long getMinIntervalSeconds() {
-    return MIN_INTERVAL_SECONDS;
-  }
 }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalJobExecutor.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalJobExecutor.java
index ca972f1cdd0..3226b3947f0 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalJobExecutor.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalJobExecutor.java
@@ -21,102 +21,19 @@ package org.apache.iotdb.commons.pipe.agent.runtime;
 
 import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
 import org.apache.iotdb.commons.concurrent.ThreadName;
-import org.apache.iotdb.commons.concurrent.WrappedRunnable;
-import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil;
 import org.apache.iotdb.commons.pipe.config.PipeConfig;
-import org.apache.iotdb.commons.utils.TestOnly;
-
-import org.apache.tsfile.utils.Pair;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.List;
-import java.util.concurrent.CopyOnWriteArrayList;
-import java.util.concurrent.Future;
-import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.TimeUnit;
 
 /**
- * Single thread to execute pipe periodical jobs on DataNode or ConfigNode. 
This is for limiting the
- * thread num on the DataNode or ConfigNode instance.
+ * The shortest scheduling cycle for these jobs is {@link
+ * PipeConfig#getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds()}, 
suitable for jobs that are
+ * NOT time-critical.
  */
-public class PipePeriodicalJobExecutor {
-
-  private static final Logger LOGGER = 
LoggerFactory.getLogger(PipePeriodicalJobExecutor.class);
-
-  private static final ScheduledExecutorService PERIODICAL_JOB_EXECUTOR =
-      IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(
-          ThreadName.PIPE_RUNTIME_PERIODICAL_JOB_EXECUTOR.getName());
-
-  private static final long MIN_INTERVAL_SECONDS =
-      
PipeConfig.getInstance().getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds();
-
-  private long rounds;
-  private Future<?> executorFuture;
-
-  // <Periodical job, Interval in rounds>
-  private final List<Pair<WrappedRunnable, Long>> periodicalJobs = new 
CopyOnWriteArrayList<>();
-
-  public void register(String id, Runnable periodicalJob, long 
intervalInSeconds) {
-    periodicalJobs.add(
-        new Pair<>(
-            new WrappedRunnable() {
-              @Override
-              public void runMayThrow() {
-                try {
-                  periodicalJob.run();
-                } catch (Exception e) {
-                  LOGGER.warn("Periodical job {} failed.", id, e);
-                }
-              }
-            },
-            Math.max(intervalInSeconds / MIN_INTERVAL_SECONDS, 1)));
-    LOGGER.info(
-        "Pipe periodical job {} is registered successfully. Interval: {} 
seconds.",
-        id,
-        Math.max(intervalInSeconds / MIN_INTERVAL_SECONDS, 1) * 
MIN_INTERVAL_SECONDS);
-  }
-
-  public synchronized void start() {
-    if (executorFuture == null) {
-      rounds = 0;
-
-      executorFuture =
-          ScheduledExecutorUtil.safelyScheduleWithFixedDelay(
-              PERIODICAL_JOB_EXECUTOR,
-              this::execute,
-              MIN_INTERVAL_SECONDS,
-              MIN_INTERVAL_SECONDS,
-              TimeUnit.SECONDS);
-      LOGGER.info("Pipe periodical job executor is started successfully.");
-    }
-  }
-
-  private void execute() {
-    ++rounds;
-
-    for (final Pair<WrappedRunnable, Long> periodicalJob : periodicalJobs) {
-      if (rounds % periodicalJob.right == 0) {
-        periodicalJob.left.run();
-      }
-    }
-  }
-
-  public synchronized void stop() {
-    if (executorFuture != null) {
-      executorFuture.cancel(false);
-      executorFuture = null;
-      LOGGER.info("Pipe periodical job executor is stopped successfully.");
-    }
-  }
-
-  @TestOnly
-  public void clear() {
-    periodicalJobs.clear();
-    LOGGER.info("All pipe periodical jobs are cleared successfully.");
-  }
+public class PipePeriodicalJobExecutor extends 
AbstractPipePeriodicalJobExecutor {
 
-  public static long getMinIntervalSeconds() {
-    return MIN_INTERVAL_SECONDS;
+  public PipePeriodicalJobExecutor() {
+    super(
+        IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(
+            ThreadName.PIPE_RUNTIME_PERIODICAL_JOB_EXECUTOR.getName()),
+        
PipeConfig.getInstance().getPipeSubtaskExecutorCronHeartbeatEventIntervalSeconds());
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/ref/PipeDataNodePhantomReferenceManager.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalPhantomReferenceCleaner.java
similarity index 54%
copy from 
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/ref/PipeDataNodePhantomReferenceManager.java
copy to 
iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalPhantomReferenceCleaner.java
index 06faed209a9..32f1549917c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/ref/PipeDataNodePhantomReferenceManager.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/runtime/PipePeriodicalPhantomReferenceCleaner.java
@@ -17,22 +17,18 @@
  * under the License.
  */
 
-package org.apache.iotdb.db.pipe.resource.ref;
+package org.apache.iotdb.commons.pipe.agent.runtime;
 
-import org.apache.iotdb.commons.pipe.config.PipeConfig;
-import org.apache.iotdb.commons.pipe.resource.ref.PipePhantomReferenceManager;
-import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
+import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
+import org.apache.iotdb.commons.concurrent.ThreadName;
 
-public class PipeDataNodePhantomReferenceManager extends 
PipePhantomReferenceManager {
+/** The shortest scheduling cycle for these jobs is 1, suitable for jobs that 
are time-critical. */
+public class PipePeriodicalPhantomReferenceCleaner extends 
AbstractPipePeriodicalJobExecutor {
 
-  public PipeDataNodePhantomReferenceManager() {
-    super();
-
-    PipeDataNodeAgent.runtime()
-        .registerPeriodicalJob(
-            "PipePhantomReferenceManager#gcHook()",
-            // NOTE: lambda CAN NOT be replaced with method reference
-            () -> super.gcHook(),
-            
PipeConfig.getInstance().getPipeEventReferenceEliminateIntervalSeconds());
+  public PipePeriodicalPhantomReferenceCleaner() {
+    super(
+        IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(
+            
ThreadName.PIPE_RUNTIME_PERIODICAL_PHANTOM_REFERENCE_CLEANER.getName()),
+        1L);
   }
 }
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/resource/ref/PipePhantomReferenceManager.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/resource/ref/PipePhantomReferenceManager.java
index 7422221c3d3..e53ff2aaa49 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/resource/ref/PipePhantomReferenceManager.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/resource/ref/PipePhantomReferenceManager.java
@@ -55,19 +55,35 @@ public abstract class PipePhantomReferenceManager {
       return;
     }
 
+    final long startTime = System.currentTimeMillis();
+
+    // limit control to avoid infinite execution
+    final int maxCount = getPhantomReferenceCount();
+    int count = 0;
+
     Reference<? extends EnrichedEvent> reference;
     try {
-      while ((reference = REFERENCE_QUEUE.remove(500)) != null) {
+      while (count < maxCount && ((reference = REFERENCE_QUEUE.remove(500)) != 
null)) {
         finalizeResource((PipeEventPhantomReference) reference);
+        count++;
       }
     } catch (final InterruptedException e) {
       // Finalize remaining references.
-      while ((reference = REFERENCE_QUEUE.poll()) != null) {
+      while (count < maxCount && ((reference = REFERENCE_QUEUE.poll()) != 
null)) {
         finalizeResource((PipeEventPhantomReference) reference);
+        count++;
       }
     } catch (final Exception e) {
       // Nowhere to really log this.
     }
+
+    if (maxCount != 0 || getPhantomReferenceCount() != 0) {
+      LOGGER.info(
+          "Clean {} pipe phantom reference(s) successfully within {} ms, 
remaining reference count: {}",
+          count,
+          System.currentTimeMillis() - startTime,
+          getPhantomReferenceCount());
+    }
   }
 
   private void finalizeResource(final PipeEventPhantomReference reference) {

Reply via email to