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

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


The following commit(s) were added to refs/heads/master by this push:
     new 36f6d6f00 [INLONG-6251][Agent] Fix the ConcurrentModification error in 
the unit test (#6252)
36f6d6f00 is described below

commit 36f6d6f00ae66ec59996857f318ef6d4ac6311a2
Author: ganfengtan <[email protected]>
AuthorDate: Sun Oct 23 14:41:28 2022 +0800

    [INLONG-6251][Agent] Fix the ConcurrentModification error in the unit test 
(#6252)
---
 .../inlong/agent/constant/AgentConstants.java      | 33 ++--------------------
 .../inlong/agent/constant/CommonConstants.java     |  2 +-
 .../apache/inlong/agent/constant/JobConstants.java |  2 +-
 .../apache/inlong/agent/core/job/JobWrapper.java   |  4 ++-
 .../apache/inlong/agent/core/task/TaskWrapper.java |  2 +-
 5 files changed, 9 insertions(+), 34 deletions(-)

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 485ce3e7a..c332d154c 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
@@ -43,88 +43,61 @@ public class AgentConstants {
     public static final String DEFAULT_AGENT_ROCKS_DB_PATH = ".rocksdb";
 
     public static final String AGENT_UNIQ_ID = "agent.uniq.id";
-    // default use local ip as uniq id for agent.
-    public static final String DEFAULT_AGENT_UNIQ_ID = AgentUtils.getLocalIp();
-
     public static final String AGENT_DB_INSTANCE_NAME = 
"agent.db.instance.name";
     public static final String DEFAULT_AGENT_DB_INSTANCE_NAME = "agent";
-
     public static final String AGENT_DB_CLASSNAME = "agent.db.classname";
     public static final String DEFAULT_AGENT_DB_CLASSNAME = 
"org.apache.inlong.agent.db.RocksDbImp";
-
     // default is empty.
     public static final String AGENT_FETCHER_CLASSNAME = 
"agent.fetcher.classname";
-
     public static final String AGENT_MESSAGE_FILTER_CLASSNAME = 
"agent.message.filter.classname";
-
     public static final String AGENT_CONF_PARENT = "agent.conf.parent";
     public static final String DEFAULT_AGENT_CONF_PARENT = "conf";
-
     public static final String AGENT_HTTP_PORT = "agent.http.port";
     public static final int DEFAULT_AGENT_HTTP_PORT = 8008;
-
     public static final String AGENT_ENABLE_HTTP = "agent.http.enable";
     public static final boolean DEFAULT_AGENT_ENABLE_HTTP = false;
-
     public static final String TRIGGER_FETCH_INTERVAL = 
"trigger.fetch.interval";
     public static final int DEFAULT_TRIGGER_FETCH_INTERVAL = 1;
-
     public static final String TRIGGER_MAX_RUNNING_NUM = 
"trigger.max.running.num";
     public static final int DEFAULT_TRIGGER_MAX_RUNNING_NUM = 4096;
-
     public static final String AGENT_FETCH_CENTER_INTERVAL_SECONDS = 
"agent.fetchCenter.interval";
     public static final int DEFAULT_AGENT_FETCH_CENTER_INTERVAL_SECONDS = 5;
-
     public static final String AGENT_TRIGGER_CHECK_INTERVAL_SECONDS = 
"agent.trigger.check.interval";
     public static final int DEFAULT_AGENT_TRIGGER_CHECK_INTERVAL_SECONDS = 1;
-
     public static final String THREAD_POOL_AWAIT_TIME = 
"thread.pool.await.time";
     // time in ms
     public static final long DEFAULT_THREAD_POOL_AWAIT_TIME = 300;
-
     public static final String JOB_MONITOR_INTERVAL = "job.monitor.interval";
     public static final int DEFAULT_JOB_MONITOR_INTERVAL = 5;
-
     public static final String JOB_FINISH_CHECK_INTERVAL = 
"job.finish.checkInterval";
     public static final long DEFAULT_JOB_FINISH_CHECK_INTERVAL = 6L;
-
     public static final String TASK_RETRY_MAX_CAPACITY = 
"task.retry.maxCapacity";
     public static final int DEFAULT_TASK_RETRY_MAX_CAPACITY = 10000;
-
     public static final String TASK_MONITOR_INTERVAL = "task.monitor.interval";
     public static final int DEFAULT_TASK_MONITOR_INTERVAL = 6;
-
     public static final String TASK_RETRY_SUBMIT_WAIT_SECONDS = 
"task.retry.submit.waitSeconds";
     public static final int DEFAULT_TASK_RETRY_SUBMIT_WAIT_SECONDS = 5;
-
     public static final String TASK_MAX_RETRY_TIME = "task.maxRetry.time";
     public static final int DEFAULT_TASK_MAX_RETRY_TIME = 3;
-
     public static final String TASK_PUSH_MAX_SECOND = "task.push.maxSecond";
     public static final int DEFAULT_TASK_PUSH_MAX_SECOND = 2;
-
     public static final String TASK_PULL_MAX_SECOND = "task.pull.maxSecond";
     public static final int DEFAULT_TASK_PULL_MAX_SECOND = 2;
-
     public static final String CHANNEL_MEMORY_CAPACITY = 
"channel.memory.capacity";
-    public static final int DEFAULT_CHANNEL_MEMORY_CAPACITY = 1000;
-
+    public static final int DEFAULT_CHANNEL_MEMORY_CAPACITY = 2000;
     public static final String TRIGGER_CHECK_INTERVAL = 
"trigger.check.interval";
     public static final int DEFAULT_TRIGGER_CHECK_INTERVAL = 2;
-
     public static final String JOB_DB_CACHE_TIME = "job.db.cache.time";
     // cache for 3 days.
     public static final long DEFAULT_JOB_DB_CACHE_TIME = 3 * 24 * 60 * 60 * 
1000;
-
     public static final String JOB_DB_CACHE_CHECK_INTERVAL = 
"job.db.cache.check.interval";
     public static final int DEFAULT_JOB_DB_CACHE_CHECK_INTERVAL = 60 * 60;
-
     public static final String JOB_NUMBER_LIMIT = "job.number.limit";
     public static final int DEFAULT_JOB_NUMBER_LIMIT = 15;
-
     public static final String AGENT_LOCAL_IP = "agent.local.ip";
     public static final String DEFAULT_LOCAL_IP = "127.0.0.1";
-
+    // default use local ip as uniq id for agent.
+    public static final String DEFAULT_AGENT_UNIQ_ID = AgentUtils.getLocalIp();
     public static final String CUSTOM_FIXED_IP = "agent.custom.fixed.ip";
 
     public static final String AGENT_CLUSTER_NAME = "agent.cluster.name";
diff --git 
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/CommonConstants.java
 
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/CommonConstants.java
index 9ecfd8ca7..6418a1252 100644
--- 
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/CommonConstants.java
+++ 
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/CommonConstants.java
@@ -77,7 +77,7 @@ public class CommonConstants {
     public static final int DEFAULT_PROXY_PACKAGE_MAX_TIMEOUT_MS = 4 * 1000;
 
     public static final String PROXY_BATCH_FLUSH_INTERVAL = 
"proxy.batch.flush.interval";
-    public static final int DEFAULT_PROXY_BATCH_FLUSH_INTERVAL = 2 * 1000;
+    public static final int DEFAULT_PROXY_BATCH_FLUSH_INTERVAL = 1000;
 
     public static final String PROXY_SENDER_MAX_TIMEOUT = 
"proxy.sender.maxTimeout";
     // max timeout in seconds.
diff --git 
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/JobConstants.java
 
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/JobConstants.java
index 87b7303ff..755ab1863 100755
--- 
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/JobConstants.java
+++ 
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/JobConstants.java
@@ -133,7 +133,7 @@ public class JobConstants extends CommonConstants {
 
     public static final String JOB_READ_WAIT_TIMEOUT = "job.file.read.wait";
 
-    public static final int DEFAULT_JOB_READ_WAIT_TIMEOUT = 5;
+    public static final int DEFAULT_JOB_READ_WAIT_TIMEOUT = 3;
 
     public static final String JOB_ID_PREFIX = "job_";
 
diff --git 
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/JobWrapper.java
 
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/JobWrapper.java
index 160c7d685..3b6ba8879 100644
--- 
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/JobWrapper.java
+++ 
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/job/JobWrapper.java
@@ -33,6 +33,7 @@ import org.slf4j.LoggerFactory;
 
 import java.util.ArrayList;
 import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
 import java.util.concurrent.TimeUnit;
 
 import static 
org.apache.inlong.agent.constant.AgentConstants.DEFAULT_JOB_VERSION;
@@ -60,7 +61,7 @@ public class JobWrapper extends AbstractStateWrapper {
         this.taskManager = manager.getTaskManager();
         this.jobManager = manager.getJobManager();
         this.job = job;
-        this.allTasks = new ArrayList<>();
+        this.allTasks = new CopyOnWriteArrayList<>();
         this.db = manager.getCommandDb();
         doChangeState(State.ACCEPTED);
     }
@@ -149,6 +150,7 @@ public class JobWrapper extends AbstractStateWrapper {
             submitAllTasks();
             checkAllTasksStateAndWait();
             cleanup();
+            LOGGER.info("job name is {}, state is {}", job.getName(), 
getCurrentState());
         } catch (Exception ex) {
             doChangeState(State.FAILED);
             LOGGER.error("error caught: {}, message: {}",
diff --git 
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskWrapper.java
 
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskWrapper.java
index 7b6c03478..6e6b8132f 100755
--- 
a/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskWrapper.java
+++ 
b/inlong-agent/agent-core/src/main/java/org/apache/inlong/agent/core/task/TaskWrapper.java
@@ -210,7 +210,7 @@ public class TaskWrapper extends AbstractStateWrapper {
             if (!isException()) {
                 doChangeState(State.SUCCEEDED);
             }
-            LOGGER.info("start to destroy task {}", task.getTaskId());
+            LOGGER.info("task state is {}, start to destroy task {}", 
getCurrentState(), task.getTaskId());
             task.destroy();
         } catch (Exception ex) {
             LOGGER.error("error while running wrapper", ex);

Reply via email to