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