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) {