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
 
 
 ###########################

Reply via email to