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/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new f1d6d1733 [INLONG-4535][Agent] Support configurable automatic exit
function when OOM happens (#4536)
f1d6d1733 is described below
commit f1d6d173355437a2b2c1dea1e3ea73eebf27bc1d
Author: xueyingzhang <[email protected]>
AuthorDate: Wed Jun 15 11:00:31 2022 +0800
[INLONG-4535][Agent] Support configurable automatic exit function when OOM
happens (#4536)
---
.../inlong/agent/common/AgentThreadFactory.java | 5 ++
.../inlong/agent/constant/AgentConstants.java | 3 +
.../org/apache/inlong/agent/utils/AgentUtils.java | 6 ++
.../ThreadUtils.java} | 43 ++++++-----
.../apache/inlong/agent/core/HeartbeatManager.java | 4 +-
.../java/org/apache/inlong/agent/core/job/Job.java | 4 +-
.../apache/inlong/agent/core/job/JobManager.java | 9 ++-
.../apache/inlong/agent/core/task/TaskManager.java | 4 +-
.../agent/core/task/TaskPositionManager.java | 4 +-
.../inlong/agent/core/trigger/TriggerManager.java | 15 ++--
.../agent/plugin/fetcher/ManagerFetcher.java | 4 +-
.../inlong/agent/plugin/sinks/ProxySink.java | 8 +-
.../sources/snapshot/BinlogSnapshotBase.java | 8 +-
.../agent/plugin/trigger/DirectoryTrigger.java | 4 +-
.../inlong/agent/plugin/trigger/PathPattern.java | 6 +-
.../apache/inlong/agent/plugin/TestOOMExit.java | 90 ++++++++++++++++++++++
inlong-agent/conf/agent.properties | 1 +
17 files changed, 183 insertions(+), 35 deletions(-)
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/common/AgentThreadFactory.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/common/AgentThreadFactory.java
index 12c87c584..0efeaa415 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/common/AgentThreadFactory.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/common/AgentThreadFactory.java
@@ -17,6 +17,8 @@
package org.apache.inlong.agent.common;
+import org.apache.inlong.agent.utils.AgentUtils;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -41,6 +43,9 @@ public class AgentThreadFactory implements ThreadFactory {
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, threadType + "-running-thread-" +
mThreadNum.getAndIncrement());
+ if (AgentUtils.enableOOMExit()) {
+ t.setUncaughtExceptionHandler(ThreadUtils::threadThrowableHandler);
+ }
LOGGER.debug("{} created", t.getName());
return t;
}
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/AgentConstants.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/AgentConstants.java
index e92f6e369..6c90f7f11 100755
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/AgentConstants.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/AgentConstants.java
@@ -191,4 +191,7 @@ public class AgentConstants {
public static final String JOB_VERSION = "job.version";
public static final Integer DEFAULT_JOB_VERSION = 1;
+ public static final String AGENT_ENABLE_OOM_EXIT = "agent.enable.oom.exit";
+ public static final boolean DEFAULT_ENABLE_OOM_EXIT = false;
+
}
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/utils/AgentUtils.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/utils/AgentUtils.java
index 9f5e7aec7..c5fb942d0 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/utils/AgentUtils.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/utils/AgentUtils.java
@@ -52,10 +52,12 @@ import java.util.concurrent.atomic.AtomicLong;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
+import static
org.apache.inlong.agent.constant.AgentConstants.AGENT_ENABLE_OOM_EXIT;
import static org.apache.inlong.agent.constant.AgentConstants.AGENT_LOCAL_IP;
import static org.apache.inlong.agent.constant.AgentConstants.AGENT_LOCAL_UUID;
import static
org.apache.inlong.agent.constant.AgentConstants.AGENT_LOCAL_UUID_OPEN;
import static
org.apache.inlong.agent.constant.AgentConstants.DEFAULT_AGENT_LOCAL_UUID_OPEN;
+import static
org.apache.inlong.agent.constant.AgentConstants.DEFAULT_ENABLE_OOM_EXIT;
import static
org.apache.inlong.agent.constant.FetcherConstants.DEFAULT_LOCAL_IP;
/**
@@ -429,4 +431,8 @@ public class AgentUtils {
return finalPath;
}
+ public static boolean enableOOMExit() {
+ return
AgentConfiguration.getAgentConf().getBoolean(AGENT_ENABLE_OOM_EXIT,
DEFAULT_ENABLE_OOM_EXIT);
+ }
+
}
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/common/AgentThreadFactory.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/utils/ThreadUtils.java
similarity index 51%
copy from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/common/AgentThreadFactory.java
copy to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/utils/ThreadUtils.java
index 12c87c584..b9829ca1f 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/common/AgentThreadFactory.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/utils/ThreadUtils.java
@@ -15,33 +15,40 @@
* limitations under the License.
*/
-package org.apache.inlong.agent.common;
+package org.apache.inlong.agent.utils;
+import org.apache.commons.lang.exception.ExceptionUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.util.concurrent.ThreadFactory;
-import java.util.concurrent.atomic.AtomicInteger;
-
/**
- * AgentThreadFactory, used for creating thread.
+ * ThreadUtils, used for handle specified throwable, such as oom, etc.
*/
-public class AgentThreadFactory implements ThreadFactory {
-
- private static final Logger LOGGER =
LoggerFactory.getLogger(AgentThreadFactory.class);
+public class ThreadUtils {
- private final AtomicInteger mThreadNum = new AtomicInteger(1);
+ private static final Logger LOGGER =
LoggerFactory.getLogger(ThreadUtils.class);
- private final String threadType;
+ public static void threadThrowableHandler(Thread t, Throwable e) {
+ if (AgentUtils.enableOOMExit()) {
+ handleOOM(t, e);
+ }
+ }
- public AgentThreadFactory(String threadType) {
- this.threadType = threadType;
+ private static void handleOOM(Thread t, Throwable e) {
+ if (ExceptionUtils.indexOfThrowable(e,
java.lang.OutOfMemoryError.class) != -1) {
+ LOGGER.error("Agent exit caused by {} OutOfMemory: ", t.getName(),
e);
+ forceShutDown();
+ }
}
- @Override
- public Thread newThread(Runnable r) {
- Thread t = new Thread(r, threadType + "-running-thread-" +
mThreadNum.getAndIncrement());
- LOGGER.debug("{} created", t.getName());
- return t;
+ private static void forceShutDown() {
+ try {
+ Runtime.getRuntime().exit(-1);
+ } catch (Throwable e) {
+ LOGGER.error("exit failed, just halt, exception: ", e);
+ Runtime.getRuntime().halt(-2);
+ }
}
-}
\ No newline at end of file
+
+}
+
diff --git
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/HeartbeatManager.java
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/HeartbeatManager.java
index 035394e19..91ee4b43c 100644
---
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/HeartbeatManager.java
+++
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/HeartbeatManager.java
@@ -24,6 +24,7 @@ import org.apache.inlong.agent.core.job.JobManager;
import org.apache.inlong.agent.core.job.JobWrapper;
import org.apache.inlong.agent.utils.AgentUtils;
import org.apache.inlong.agent.utils.HttpManager;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.apache.inlong.common.pojo.agent.TaskSnapshotMessage;
import org.apache.inlong.common.pojo.agent.TaskSnapshotRequest;
import org.slf4j.Logger;
@@ -138,8 +139,9 @@ public class HeartbeatManager extends AbstractDaemon {
int heartbeatInterval =
conf.getInt(AGENT_HEARTBEAT_INTERVAL,
DEFAULT_AGENT_HEARTBEAT_INTERVAL);
TimeUnit.SECONDS.sleep(heartbeatInterval);
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("error caught", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
ex);
}
}
};
diff --git
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/Job.java
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/Job.java
index 42b127636..9b2cf11e2 100644
---
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/Job.java
+++
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/Job.java
@@ -24,6 +24,7 @@ import org.apache.inlong.agent.plugin.Channel;
import org.apache.inlong.agent.plugin.Reader;
import org.apache.inlong.agent.plugin.Sink;
import org.apache.inlong.agent.plugin.Source;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -97,8 +98,9 @@ public class Job {
String taskId = String.format("%s_%d", jobInstanceId, index++);
taskList.add(new Task(taskId, reader, writer, channel,
getJobConf()));
}
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("create task failed", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(), ex);
throw new RuntimeException(ex);
}
return taskList;
diff --git
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/JobManager.java
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/JobManager.java
index 315594121..67fb4149c 100644
---
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/JobManager.java
+++
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/JobManager.java
@@ -27,6 +27,7 @@ import org.apache.inlong.agent.db.JobProfileDb;
import org.apache.inlong.agent.db.StateSearchKey;
import org.apache.inlong.agent.utils.AgentUtils;
import org.apache.inlong.agent.utils.ConfigUtil;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -121,6 +122,8 @@ public class JobManager extends AbstractDaemon {
} catch (Exception rje) {
LOGGER.debug("reject job {}", job.getJobInstanceId(), rje);
pendingJobs.putIfAbsent(job.getJobInstanceId(), job);
+ } catch (Throwable t) {
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(), t);
}
}
@@ -233,8 +236,9 @@ public class JobManager extends AbstractDaemon {
}
}
TimeUnit.SECONDS.sleep(monitorInterval);
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("error caught", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
ex);
}
}
};
@@ -253,8 +257,9 @@ public class JobManager extends AbstractDaemon {
}
try {
TimeUnit.SECONDS.sleep(jobDbCacheCheckInterval);
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("sleep error caught", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
ex);
}
}
};
diff --git
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskManager.java
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskManager.java
index 39deb0c2f..918bb4483 100755
---
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskManager.java
+++
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskManager.java
@@ -24,6 +24,7 @@ import org.apache.inlong.agent.constant.AgentConstants;
import org.apache.inlong.agent.core.AgentManager;
import org.apache.inlong.agent.utils.AgentUtils;
import org.apache.inlong.agent.utils.ConfigUtil;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -236,8 +237,9 @@ public class TaskManager extends AbstractDaemon {
}
}
TimeUnit.SECONDS.sleep(monitorInterval);
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("Exception caught", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
ex);
}
}
};
diff --git
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskPositionManager.java
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskPositionManager.java
index b3ca813bf..b0094abf8 100644
---
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskPositionManager.java
+++
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskPositionManager.java
@@ -22,6 +22,7 @@ import org.apache.inlong.agent.conf.AgentConfiguration;
import org.apache.inlong.agent.conf.JobProfile;
import org.apache.inlong.agent.core.AgentManager;
import org.apache.inlong.agent.db.JobProfileDb;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -100,8 +101,9 @@ public class TaskPositionManager extends AbstractDaemon {
int flushTime = conf.getInt(AGENT_HEARTBEAT_INTERVAL,
DEFAULT_AGENT_FETCHER_INTERVAL);
TimeUnit.SECONDS.sleep(flushTime);
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("error caught", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
ex);
}
}
};
diff --git
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/trigger/TriggerManager.java
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/trigger/TriggerManager.java
index c4ecf512f..4ede279ea 100755
---
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/trigger/TriggerManager.java
+++
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/trigger/TriggerManager.java
@@ -28,6 +28,7 @@ import org.apache.inlong.agent.core.AgentManager;
import org.apache.inlong.agent.core.job.JobWrapper;
import org.apache.inlong.agent.db.TriggerProfileDb;
import org.apache.inlong.agent.plugin.Trigger;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -84,8 +85,9 @@ public class TriggerManager extends AbstractDaemon {
triggerMap.put(triggerId, trigger);
trigger.init(triggerProfile);
trigger.run();
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("exception caught", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(), ex);
return false;
}
return true;
@@ -140,8 +142,10 @@ public class TriggerManager extends AbstractDaemon {
}
});
TimeUnit.SECONDS.sleep(triggerFetchInterval);
- } catch (Exception ignored) {
- LOGGER.info("ignored Exception ", ignored);
+ } catch (Throwable e) {
+ LOGGER.info("ignored Exception ", e);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
e);
+
}
}
@@ -180,8 +184,9 @@ public class TriggerManager extends AbstractDaemon {
}
});
TimeUnit.MINUTES.sleep(JOB_CHECK_INTERVAL);
- } catch (Exception ignored) {
- LOGGER.info("ignored Exception ", ignored);
+ } catch (Throwable e) {
+ LOGGER.info("ignored Exception ", e);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
e);
}
}
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/fetcher/ManagerFetcher.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/fetcher/ManagerFetcher.java
index b6a06e56b..3b0d4f840 100755
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/fetcher/ManagerFetcher.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/fetcher/ManagerFetcher.java
@@ -37,6 +37,7 @@ import org.apache.inlong.agent.pojo.DbCollectorTaskRequestDto;
import org.apache.inlong.agent.pojo.DbCollectorTaskResult;
import org.apache.inlong.agent.utils.AgentUtils;
import org.apache.inlong.agent.utils.HttpManager;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.apache.inlong.common.db.CommandEntity;
import org.apache.inlong.common.enums.ManagerOpEnum;
import org.apache.inlong.common.enums.PullJobTypeEnum;
@@ -489,8 +490,9 @@ public class ManagerFetcher extends AbstractDaemon
implements ProfileFetcher {
// fetch db collector task from manager
fetchDbCollectTask();
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.warn("exception caught", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
ex);
}
}
};
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sinks/ProxySink.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sinks/ProxySink.java
index 235f9730e..97fa7d3d8 100755
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sinks/ProxySink.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sinks/ProxySink.java
@@ -28,6 +28,7 @@ import org.apache.inlong.agent.plugin.Message;
import org.apache.inlong.agent.plugin.MessageFilter;
import org.apache.inlong.agent.plugin.message.PackProxyMessage;
import org.apache.inlong.agent.utils.AgentUtils;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -120,6 +121,8 @@ public class ProxySink extends AbstractSink {
}
} catch (Exception e) {
LOGGER.error("write message to Proxy sink error", e);
+ } catch (Throwable t) {
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(), t);
}
}
@@ -170,6 +173,8 @@ public class ProxySink extends AbstractSink {
AgentUtils.silenceSleepInMs(batchFlushInterval);
} catch (Exception ex) {
LOGGER.error("error caught", ex);
+ } catch (Throwable t) {
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
t);
}
}
};
@@ -196,8 +201,9 @@ public class ProxySink extends AbstractSink {
senderManager = new SenderManager(jobConf, inlongGroupId, sourceName);
try {
senderManager.addMessageSender();
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("error while init sender for group id {}",
inlongGroupId);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(), ex);
throw new IllegalStateException(ex);
}
}
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/snapshot/BinlogSnapshotBase.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/snapshot/BinlogSnapshotBase.java
index a01deb19d..238875ed3 100644
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/snapshot/BinlogSnapshotBase.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/snapshot/BinlogSnapshotBase.java
@@ -17,6 +17,7 @@
package org.apache.inlong.agent.plugin.sources.snapshot;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -76,8 +77,9 @@ public class BinlogSnapshotBase implements SnapshotBase {
offset = outputStream.toByteArray();
inputStream.close();
outputStream.close();
- } catch (Exception ex) {
+ } catch (Throwable ex) {
log.error("load binlog offset error", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(), ex);
}
}
@@ -90,8 +92,10 @@ public class BinlogSnapshotBase implements SnapshotBase {
offset = bytes;
try (OutputStream output = new FileOutputStream(file)) {
output.write(bytes);
- } catch (Exception e) {
+ } catch (Throwable e) {
log.error("save offset to file error", e);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(), e);
+
}
}
}
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/trigger/DirectoryTrigger.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/trigger/DirectoryTrigger.java
index 743eb0bd9..222969dd3 100644
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/trigger/DirectoryTrigger.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/trigger/DirectoryTrigger.java
@@ -24,6 +24,7 @@ import org.apache.inlong.agent.constant.AgentConstants;
import org.apache.inlong.agent.constant.JobConstants;
import org.apache.inlong.agent.plugin.Trigger;
import org.apache.inlong.agent.plugin.utils.PluginUtils;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -200,8 +201,9 @@ public class DirectoryTrigger extends AbstractDaemon
implements Trigger {
watchKeys.addAll(tmpWatchers);
watchKeys.removeAll(tmpDeletedWatchers);
});
- } catch (Exception ex) {
+ } catch (Throwable ex) {
LOGGER.error("error caught", ex);
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(),
ex);
}
}
};
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/trigger/PathPattern.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/trigger/PathPattern.java
index 7a068f46b..dda014e32 100644
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/trigger/PathPattern.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/trigger/PathPattern.java
@@ -20,6 +20,7 @@ package org.apache.inlong.agent.plugin.trigger;
import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.builder.HashCodeBuilder;
import org.apache.inlong.agent.plugin.filter.DateFormatRegex;
+import org.apache.inlong.agent.utils.ThreadUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -98,7 +99,10 @@ public class PathPattern {
}
});
} catch (Exception e) {
- LOGGER.error("error caught", e);
+ LOGGER.error("exception caught", e);
+ } catch (Throwable t) {
+ ThreadUtils.threadThrowableHandler(Thread.currentThread(), t);
+
}
}
}
diff --git
a/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/TestOOMExit.java
b/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/TestOOMExit.java
new file mode 100644
index 000000000..4bc44a663
--- /dev/null
+++
b/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/TestOOMExit.java
@@ -0,0 +1,90 @@
+/*
+ * 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.agent.plugin;
+
+import org.apache.inlong.agent.common.AbstractDaemon;
+import org.apache.inlong.agent.conf.AgentConfiguration;
+import org.apache.inlong.agent.utils.ThreadUtils;
+import org.junit.BeforeClass;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.powermock.api.mockito.PowerMockito;
+import org.powermock.core.classloader.annotations.PowerMockIgnore;
+import org.powermock.core.classloader.annotations.PrepareForTest;
+import org.powermock.modules.junit4.PowerMockRunner;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.concurrent.TimeUnit;
+
+import static
org.apache.inlong.agent.constant.AgentConstants.AGENT_ENABLE_OOM_EXIT;
+
+@RunWith(PowerMockRunner.class)
+@PrepareForTest(ThreadUtils.class)
+@PowerMockIgnore({"javax.management.*"})
+public class TestOOMExit {
+
+ @BeforeClass
+ public static void setup() throws Exception {
+ PowerMockito.spy(ThreadUtils.class);
+ PowerMockito.doNothing().when(ThreadUtils.class, "forceShutDown");
+ }
+
+ @Test
+ public void testOOM() {
+ MockJobManager jobManager = new MockJobManager();
+ AgentConfiguration conf = AgentConfiguration.getAgentConf();
+ conf.setBoolean(AGENT_ENABLE_OOM_EXIT, true);
+ jobManager.start();
+ jobManager.join();
+ }
+
+ static class MockJobManager extends AbstractDaemon {
+ private static final Logger LOGGER =
LoggerFactory.getLogger(MockJobManager.class);
+
+ @Override
+ public void start() {
+ submitWorker(throwOOMThread());
+ }
+
+ @Override
+ public void stop() throws Exception {
+
+ }
+
+ public Runnable throwOOMThread() {
+ return () -> {
+ int i = 0;
+ while (i < 5) {
+ try {
+ LOGGER.info("throw OOM thread: " + i);
+ TimeUnit.SECONDS.sleep(1);
+ i++;
+ if (i == 3) {
+ LOGGER.info("throw OOM");
+ throw new OutOfMemoryError();
+ }
+ } catch (Throwable ex) {
+
ThreadUtils.threadThrowableHandler(Thread.currentThread(), ex);
+ }
+ }
+ };
+ }
+ }
+
+}
diff --git a/inlong-agent/conf/agent.properties
b/inlong-agent/conf/agent.properties
index 4c353f776..2fd76acf2 100755
--- a/inlong-agent/conf/agent.properties
+++ b/inlong-agent/conf/agent.properties
@@ -53,6 +53,7 @@ thread.pool.await.time=30
agent.local.ip=127.0.0.1
agent.local.uuid=
agent.local.uuid.open=false
+agent.enable.oom.exit=false
###########################