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/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new d6f03f3c17 [INLONG-9955][Agent] Rename Job to Task for agent (#9956)
d6f03f3c17 is described below
commit d6f03f3c177ac6b3036c9e9ea0fe853fb441fb23
Author: haifxu <[email protected]>
AuthorDate: Wed Apr 10 17:00:31 2024 +0800
[INLONG-9955][Agent] Rename Job to Task for agent (#9956)
---
.../apache/inlong/agent/conf/InstanceProfile.java | 8 +-
.../inlong/agent/constant/TaskConstants.java | 214 ++++-------
.../org/apache/inlong/agent/db/JobProfileDb.java | 253 -------------
.../agent/pojo/{BinlogJob.java => BinlogTask.java} | 4 +-
.../agent/pojo/{KafkaJob.java => KafkaTask.java} | 4 +-
.../agent/pojo/{MongoJob.java => MongoTask.java} | 6 +-
.../agent/pojo/{MqttJob.java => MqttTask.java} | 6 +-
.../agent/pojo/{OracleJob.java => OracleTask.java} | 10 +-
.../{PostgreSQLJob.java => PostgreSQLTask.java} | 6 +-
.../agent/pojo/{RedisJob.java => RedisTask.java} | 4 +-
.../pojo/{SqlServerJob.java => SqlServerTask.java} | 10 +-
.../apache/inlong/agent/pojo/TaskProfileDto.java | 408 +++++++++++----------
.../plugin/sinks/filecollect/SenderManager.java | 6 +-
.../inlong/agent/plugin/sources/KafkaSource.java | 10 +-
.../inlong/agent/plugin/sources/LogFileSource.java | 10 +-
.../agent/plugin/sources/reader/MongoDBReader.java | 122 +++---
.../inlong/agent/plugin/utils/MetaDataUtils.java | 16 +-
.../inlong/agent/plugin/utils/PluginUtils.java | 10 +-
.../inlong/agent/plugin/utils/RocksDBUtils.java | 6 -
.../apache/inlong/agent/plugin/sinks/MockSink.java | 4 +-
20 files changed, 399 insertions(+), 718 deletions(-)
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/conf/InstanceProfile.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/conf/InstanceProfile.java
index 06d783c953..e3ed73495f 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/conf/InstanceProfile.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/conf/InstanceProfile.java
@@ -37,8 +37,8 @@ import static
org.apache.inlong.agent.constant.CommonConstants.DEFAULT_PROXY_INL
import static
org.apache.inlong.agent.constant.CommonConstants.PROXY_INLONG_GROUP_ID;
import static
org.apache.inlong.agent.constant.CommonConstants.PROXY_INLONG_STREAM_ID;
import static org.apache.inlong.agent.constant.TaskConstants.INSTANCE_STATE;
-import static org.apache.inlong.agent.constant.TaskConstants.JOB_MQ_ClUSTERS;
-import static org.apache.inlong.agent.constant.TaskConstants.JOB_MQ_TOPIC;
+import static org.apache.inlong.agent.constant.TaskConstants.TASK_MQ_CLUSTERS;
+import static org.apache.inlong.agent.constant.TaskConstants.TASK_MQ_TOPIC;
import static org.apache.inlong.agent.constant.TaskConstants.TASK_RETRY;
/**
@@ -127,7 +127,7 @@ public class InstanceProfile extends AbstractConfiguration
implements Comparable
*/
public List<MQClusterInfo> getMqClusters() {
List<MQClusterInfo> result = null;
- String mqClusterStr = get(JOB_MQ_ClUSTERS);
+ String mqClusterStr = get(TASK_MQ_CLUSTERS);
if (StringUtils.isNotBlank(mqClusterStr)) {
result = GSON.fromJson(mqClusterStr, new
TypeToken<List<MQClusterInfo>>() {
}.getType());
@@ -140,7 +140,7 @@ public class InstanceProfile extends AbstractConfiguration
implements Comparable
*/
public DataProxyTopicInfo getMqTopic() {
DataProxyTopicInfo result = null;
- String topicStr = get(JOB_MQ_TOPIC);
+ String topicStr = get(TASK_MQ_TOPIC);
if (StringUtils.isNotBlank(topicStr)) {
result = GSON.fromJson(topicStr, DataProxyTopicInfo.class);
}
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/TaskConstants.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/TaskConstants.java
index 138e5ee78f..23d6fe8dc2 100755
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/TaskConstants.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/constant/TaskConstants.java
@@ -18,7 +18,7 @@
package org.apache.inlong.agent.constant;
/**
- * Basic config for a single job
+ * Basic config for a single task
*/
public class TaskConstants extends CommonConstants {
@@ -29,57 +29,41 @@ public class TaskConstants extends CommonConstants {
public static final String JOB_INSTANCE_ID = "job.instance.id";
public static final String INSTANCE_CREATE_TIME = "instance.createTime";
public static final String INSTANCE_MODIFY_TIME = "instance.modifyTime";
- public static final String JOB_IP = "job.ip";
- public static final String JOB_RETRY = "job.retry";
- public static final String JOB_UUID = "job.uuid";
public static final String TASK_GROUP_ID = "task.groupId";
public static final String TASK_STREAM_ID = "task.streamId";
public static final String RESTORE_FROM_DB = "task.restoreFromDB";
public static final String TASK_SOURCE = "task.source";
- public static final String JOB_SOURCE_TYPE = "job.sourceType";
public static final String TASK_CHANNEL = "task.channel";
- public static final String JOB_NAME = "job.name";
- public static final String JOB_LINE_FILTER_PATTERN = "job.pattern";
- public static final String DEFAULT_JOB_NAME = "default";
- public static final String JOB_DESCRIPTION = "job.description";
- public static final String DEFAULT_JOB_DESCRIPTION = "default job
description";
- public static final String DEFAULT_JOB_LINE_FILTER = "";
public static final String TASK_CLASS = "task.taskClass";
public static final String INSTANCE_CLASS = "task.instance.class";
- public static final String JOB_FILE_TRIGGER = "job.fileTask.trigger";
+ public static final String TASK_FILE_TRIGGER = "task.fileTask.trigger";
// sink config
public static final String TASK_SINK = "task.sink";
- public static final String JOB_PROXY_SEND = "job.proxySend";
- public static final boolean DEFAULT_JOB_PROXY_SEND = false;
- public static final String JOB_MQ_ClUSTERS = "job.mqClusters";
- public static final String JOB_MQ_TOPIC = "job.topicInfo";
+ public static final String TASK_PROXY_SEND = "task.proxySend";
+ public static final boolean DEFAULT_TASK_PROXY_SEND = false;
+ public static final String TASK_MQ_CLUSTERS = "task.mqClusters";
+ public static final String TASK_MQ_TOPIC = "task.topicInfo";
public static final String OFFSET = "offset";
public static final Long DEFAULT_OFFSET = -1L;
public static final String INODE_INFO = "inodeInfo";
- // File job
+ // File task
public static final String TASK_DIR_FILTER_PATTERN =
"task.fileTask.dir.pattern"; // deprecated
public static final String FILE_DIR_FILTER_PATTERNS =
"task.fileTask.dir.patterns";
public static final String TASK_FILE_TIME_OFFSET =
"task.fileTask.timeOffset";
public static final String TASK_TIME_ZONE = "task.timeZone";
- public static final String TASK_FILE_MAX_WAIT =
"task.fileTask.file.max.wait";
public static final String TASK_CYCLE_UNIT = "task.cycleUnit";
public static final String FILE_TASK_CYCLE_UNIT =
"task.fileTask.cycleUnit";
- public static final String TASK_FILE_TRIGGER_TYPE =
"task.fileTask.collectType";
- public static final String JOB_FILE_LINE_END_PATTERN =
"job.fileTask.line.endPattern";
- public static final String JOB_FILE_CONTENT_COLLECT_TYPE =
"job.fileTask.contentCollectType";
- public static final String JOB_FILE_META_ENV_LIST = "job.fileTask.envList";
- public static final String JOB_FILE_META_FILTER_BY_LABELS =
"job.fileTask.filterMetaByLabels";
- public static final String JOB_FILE_PROPERTIES = "job.fileTask.properties";
+ public static final String TASK_FILE_CONTENT_COLLECT_TYPE =
"task.fileTask.contentCollectType";
+ public static final String TASK_FILE_META_ENV_LIST =
"task.fileTask.envList";
+ public static final String TASK_FILE_META_FILTER_BY_LABELS =
"task.fileTask.filterMetaByLabels";
+ public static final String TASK_FILE_PROPERTIES =
"task.fileTask.properties";
public static final String SOURCE_DATA_CONTENT_STYLE =
"task.fileTask.dataContentStyle";
public static final String SOURCE_DATA_SEPARATOR =
"task.fileTask.dataSeparator";
- public static final String JOB_FILE_MONITOR_INTERVAL =
"job.fileTask.monitorInterval";
- public static final String JOB_FILE_MONITOR_STATUS =
"job.fileTask.monitorStatus";
- public static final String JOB_FILE_MONITOR_EXPIRE =
"job.fileTask.monitorExpire";
public static final String TASK_RETRY = "task.fileTask.retry";
public static final String TASK_START_TIME = "task.fileTask.startTime";
public static final String TASK_END_TIME = "task.fileTask.endTime";
@@ -89,33 +73,39 @@ public class TaskConstants extends CommonConstants {
public static final String DEFAULT_FILE_SOURCE_EXTEND_CLASS =
"org.apache.inlong.agent.plugin.sources.file.extend.ExtendedHandler";
- // Binlog job
- public static final String JOB_DATABASE_USER = "job.binlogJob.user";
- public static final String JOB_DATABASE_PASSWORD =
"job.binlogJob.password";
- public static final String JOB_DATABASE_HOSTNAME =
"job.binlogJob.hostname";
- public static final String JOB_TABLE_WHITELIST =
"job.binlogJob.tableWhiteList";
- public static final String JOB_DATABASE_WHITELIST =
"job.binlogJob.databaseWhiteList";
- public static final String JOB_DATABASE_OFFSETS = "job.binlogJob.offsets";
- public static final String JOB_DATABASE_OFFSET_FILENAME =
"job.binlogJob.offset.filename";
-
- public static final String JOB_DATABASE_SERVER_TIME_ZONE =
"job.binlogJob.serverTimezone";
- public static final String JOB_DATABASE_STORE_OFFSET_INTERVAL_MS =
"job.binlogJob.offset.intervalMs";
-
- public static final String JOB_DATABASE_STORE_HISTORY_FILENAME =
"job.binlogJob.history.filename";
- public static final String JOB_DATABASE_INCLUDE_SCHEMA_CHANGES =
"job.binlogJob.schema";
- public static final String JOB_DATABASE_SNAPSHOT_MODE =
"job.binlogJob.snapshot.mode";
- public static final String JOB_DATABASE_HISTORY_MONITOR_DDL =
"job.binlogJob.ddl";
- public static final String JOB_DATABASE_PORT = "job.binlogJob.port";
-
- // Kafka job
- public static final String TASK_KAFKA_TOPIC = "task.kafkaJob.topic";
- public static final String TASK_KAFKA_BOOTSTRAP_SERVERS =
"task.kafkaJob.bootstrap.servers";
- public static final String TASK_KAFKA_GROUP_ID = "task.kafkaJob.group.id";
- public static final String JOB_KAFKA_RECORD_SPEED_LIMIT =
"job.kafkaJob.recordSpeed.limit";
- public static final String JOB_KAFKA_BYTE_SPEED_LIMIT =
"job.kafkaJob.byteSpeed.limit";
- public static final String TASK_KAFKA_OFFSET =
"task.kafkaJob.partition.offset";
- public static final String JOB_KAFKA_READ_TIMEOUT =
"job.kafkaJob.read.timeout";
- public static final String TASK_KAFKA_AUTO_COMMIT_OFFSET_RESET =
"task.kafkaJob.autoOffsetReset";
+ // Binlog task
+ public static final String TASK_DATABASE_USER = "task.binlogTask.user";
+ public static final String TASK_DATABASE_PASSWORD =
"task.binlogTask.password";
+ public static final String TASK_DATABASE_HOSTNAME =
"task.binlogTask.hostname";
+ public static final String TASK_TABLE_WHITELIST =
"task.binlogTask.tableWhiteList";
+ public static final String TASK_DATABASE_WHITELIST =
"task.binlogTask.databaseWhiteList";
+ public static final String TASK_DATABASE_OFFSETS =
"task.binlogTask.offsets";
+ public static final String TASK_DATABASE_OFFSET_FILENAME =
"task.binlogTask.offset.filename";
+
+ public static final String TASK_DATABASE_SERVER_TIME_ZONE =
"task.binlogTask.serverTimezone";
+ public static final String TASK_DATABASE_STORE_OFFSET_INTERVAL_MS =
"task.binlogTask.offset.intervalMs";
+
+ public static final String TASK_DATABASE_STORE_HISTORY_FILENAME =
"task.binlogTask.history.filename";
+ public static final String TASK_DATABASE_INCLUDE_SCHEMA_CHANGES =
"task.binlogTask.schema";
+ public static final String TASK_DATABASE_SNAPSHOT_MODE =
"task.binlogTask.snapshot.mode";
+ public static final String TASK_DATABASE_HISTORY_MONITOR_DDL =
"task.binlogTask.ddl";
+ public static final String TASK_DATABASE_PORT = "task.binlogTask.port";
+
+ // Kafka task
+ public static final String TASK_KAFKA_TOPIC = "task.kafkaTask.topic";
+ public static final String TASK_KAFKA_BOOTSTRAP_SERVERS =
"task.kafkaTask.bootstrap.servers";
+ public static final String TASK_KAFKA_OFFSET =
"task.kafkaTask.partition.offset";
+ public static final String TASK_KAFKA_AUTO_COMMIT_OFFSET_RESET =
"task.kafkaTask.autoOffsetReset";
+
+ /**
+ * delimiter to split offset for different task
+ */
+ public static final String TASK_KAFKA_OFFSET_DELIMITER = "_";
+
+ /**
+ * delimiter to split all partition offset for all kafka tasks
+ */
+ public static final String TASK_KAFKA_PARTITION_OFFSET_DELIMITER = "#";
// Pulsar task
public static final String TASK_PULSAR_TENANT = "task.pulsarTask.tenant";
@@ -127,58 +117,34 @@ public class TaskConstants extends CommonConstants {
public static final String TASK_PULSAR_SUBSCRIPTION_POSITION =
"task.pulsarTask.subscriptionPosition";
public static final String TASK_PULSAR_RESET_TIME =
"task.pulsarTask.resetTime";
- public static final String JOB_MONGO_HOSTS = "job.mongoJob.hosts";
- public static final String JOB_MONGO_USER = "job.mongoJob.user";
- public static final String JOB_MONGO_PASSWORD = "job.mongoJob.password";
- public static final String JOB_MONGO_DATABASE_INCLUDE_LIST =
"job.mongoJob.databaseIncludeList";
- public static final String JOB_MONGO_DATABASE_EXCLUDE_LIST =
"job.mongoJob.databaseExcludeList";
- public static final String JOB_MONGO_COLLECTION_INCLUDE_LIST =
"job.mongoJob.collectionIncludeList";
- public static final String JOB_MONGO_COLLECTION_EXCLUDE_LIST =
"job.mongoJob.collectionExcludeList";
- public static final String JOB_MONGO_FIELD_EXCLUDE_LIST =
"job.mongoJob.fieldExcludeList";
- public static final String JOB_MONGO_SNAPSHOT_MODE =
"job.mongoJob.snapshotMode";
- public static final String JOB_MONGO_CAPTURE_MODE =
"job.mongoJob.captureMode";
- public static final String JOB_MONGO_QUEUE_SIZE = "job.mongoJob.queueSize";
- public static final String JOB_MONGO_STORE_HISTORY_FILENAME =
"job.mongoJob.history.filename";
- public static final String JOB_MONGO_OFFSET_SPECIFIC_OFFSET_FILE =
"job.mongoJob.offset.specificOffsetFile";
- public static final String JOB_MONGO_OFFSET_SPECIFIC_OFFSET_POS =
"job.mongoJob.offset.specificOffsetPos";
- public static final String JOB_MONGO_OFFSETS = "job.mongoJob.offsets";
- public static final String JOB_MONGO_CONNECT_TIMEOUT_MS =
"job.mongoJob.connectTimeoutInMs";
- public static final String JOB_MONGO_CURSOR_MAX_AWAIT =
"job.mongoJob.cursorMaxAwaitTimeInMs";
- public static final String JOB_MONGO_SOCKET_TIMEOUT =
"job.mongoJob.socketTimeoutInMs";
- public static final String JOB_MONGO_SELECTION_TIMEOUT =
"job.mongoJob.selectionTimeoutInMs";
- public static final String JOB_MONGO_FIELD_RENAMES =
"job.mongoJob.fieldRenames";
- public static final String JOB_MONGO_MEMBERS_DISCOVER =
"job.mongoJob.membersAutoDiscover";
- public static final String JOB_MONGO_CONNECT_MAX_ATTEMPTS =
"job.mongoJob.connectMaxAttempts";
- public static final String JOB_MONGO_BACKOFF_MAX_DELAY =
"job.mongoJob.connectBackoffMaxDelayInMs";
- public static final String JOB_MONGO_BACKOFF_INITIAL_DELAY =
"job.mongoJob.connectBackoffInitialDelayInMs";
- public static final String JOB_MONGO_INITIAL_SYNC_MAX_THREADS =
"job.mongoJob.initialSyncMaxThreads";
- public static final String JOB_MONGO_SSL_INVALID_HOSTNAME_ALLOWED =
"job.mongoJob.sslInvalidHostnameAllowed";
- public static final String JOB_MONGO_SSL_ENABLE =
"job.mongoJob.sslEnabled";
- public static final String JOB_MONGO_POLL_INTERVAL =
"job.mongoJob.pollIntervalInMs";
-
- public static final Long JOB_KAFKA_DEFAULT_OFFSET = 0L;
-
- // job type, delete/add
- public static final String JOB_TYPE = "job.type";
-
- public static final String JOB_CHECKPOINT = "job.checkpoint";
-
- public static final String DEFAULT_JOB_FILE_TIME_OFFSET = "0d";
-
- // time in min
- public static final int DEFAULT_JOB_FILE_MAX_WAIT_TIME = 1;
-
- public static final String JOB_READ_WAIT_TIMEOUT = "job.file.read.wait";
-
- public static final String JOB_ID_PREFIX = "job_";
-
- public static final String SQL_JOB_ID = "sql_job_id";
-
- public static final String JOB_STORE_TIME = "job.store.time";
-
- public static final String JOB_OP = "job.op";
-
- public static final String JOB_STATE = "job.state";
+ public static final String TASK_MONGO_HOSTS = "task.mongoTask.hosts";
+ public static final String TASK_MONGO_USER = "task.mongoTask.user";
+ public static final String TASK_MONGO_PASSWORD = "task.mongoTask.password";
+ public static final String TASK_MONGO_DATABASE_INCLUDE_LIST =
"task.mongoTask.databaseIncludeList";
+ public static final String TASK_MONGO_DATABASE_EXCLUDE_LIST =
"task.mongoTask.databaseExcludeList";
+ public static final String TASK_MONGO_COLLECTION_INCLUDE_LIST =
"task.mongoTask.collectionIncludeList";
+ public static final String TASK_MONGO_COLLECTION_EXCLUDE_LIST =
"task.mongoTask.collectionExcludeList";
+ public static final String TASK_MONGO_FIELD_EXCLUDE_LIST =
"task.mongoTask.fieldExcludeList";
+ public static final String TASK_MONGO_SNAPSHOT_MODE =
"task.mongoTask.snapshotMode";
+ public static final String TASK_MONGO_CAPTURE_MODE =
"task.mongoTask.captureMode";
+ public static final String TASK_MONGO_QUEUE_SIZE =
"task.mongoTask.queueSize";
+ public static final String TASK_MONGO_STORE_HISTORY_FILENAME =
"task.mongoTask.history.filename";
+ public static final String TASK_MONGO_OFFSET_SPECIFIC_OFFSET_FILE =
"task.mongoTask.offset.specificOffsetFile";
+ public static final String TASK_MONGO_OFFSET_SPECIFIC_OFFSET_POS =
"task.mongoTask.offset.specificOffsetPos";
+ public static final String TASK_MONGO_OFFSETS = "task.mongoTask.offsets";
+ public static final String TASK_MONGO_CONNECT_TIMEOUT_MS =
"task.mongoTask.connectTimeoutInMs";
+ public static final String TASK_MONGO_CURSOR_MAX_AWAIT =
"task.mongoTask.cursorMaxAwaitTimeInMs";
+ public static final String TASK_MONGO_SOCKET_TIMEOUT =
"task.mongoTask.socketTimeoutInMs";
+ public static final String TASK_MONGO_SELECTION_TIMEOUT =
"task.mongoTask.selectionTimeoutInMs";
+ public static final String TASK_MONGO_FIELD_RENAMES =
"task.mongoTask.fieldRenames";
+ public static final String TASK_MONGO_MEMBERS_DISCOVER =
"task.mongoTask.membersAutoDiscover";
+ public static final String TASK_MONGO_CONNECT_MAX_ATTEMPTS =
"task.mongoTask.connectMaxAttempts";
+ public static final String TASK_MONGO_BACKOFF_MAX_DELAY =
"task.mongoTask.connectBackoffMaxDelayInMs";
+ public static final String TASK_MONGO_BACKOFF_INITIAL_DELAY =
"task.mongoTask.connectBackoffInitialDelayInMs";
+ public static final String TASK_MONGO_INITIAL_SYNC_MAX_THREADS =
"task.mongoTask.initialSyncMaxThreads";
+ public static final String TASK_MONGO_SSL_INVALID_HOSTNAME_ALLOWED =
"task.mongoTask.sslInvalidHostnameAllowed";
+ public static final String TASK_MONGO_SSL_ENABLE =
"task.mongoTask.sslEnabled";
+ public static final String TASK_MONGO_POLL_INTERVAL =
"task.mongoTask.pollIntervalInMs";
public static final String TASK_STATE = "task.state";
@@ -188,54 +154,20 @@ public class TaskConstants extends CommonConstants {
public static final String LAST_UPDATE_TIME = "lastUpdateTime";
- public static final String TRIGGER_ONLY_ONE_JOB = "job.standalone"; //
TODO:delete it
-
- // field splitter
- public static final String JOB_FIELD_SPLITTER = "job.splitter";
-
- // job delivery time
- public static final String JOB_DELIVERY_TIME = "job.deliveryTime";
-
// data time reading file
public static final String SOURCE_DATA_TIME = "source.dataTime";
// data time for sink
public static final String SINK_DATA_TIME = "sink.dataTime";
- // job of the number of seconds to wait before starting the task
- public static final String JOB_TASK_BEGIN_WAIT_SECONDS =
"job.taskWaitSeconds";
-
/**
* when job is retried, the retry time should be provided
*/
- public static final String JOB_RETRY_TIME = "job.retryTime";
-
- /**
- * delimiter to split offset for different task
- */
- public static final String JOB_OFFSET_DELIMITER = "_";
-
- /**
- * delimiter to split all partition offset for all kafka tasks
- */
- public static final String JOB_KAFKA_PARTITION_OFFSET_DELIMITER = "#";
+ public static final String TASK_RETRY_TIME = "task.retryTime";
/**
* sync send data when sending to DataProxy
*/
public static final int SYNC_SEND_OPEN = 1;
- public static final String INTERVAL_MILLISECONDS = "1000";
-
- /**
- * monitor switch, 1 true and 0 false
- */
- public static final String JOB_FILE_MONITOR_DEFAULT_STATUS = "1";
-
- /**
- * monitor expire time and the time in milliseconds.
- * default value is -1 and stand for not expire time.
- */
- public static final String JOB_FILE_MONITOR_DEFAULT_EXPIRE = "-1";
-
}
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/db/JobProfileDb.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/db/JobProfileDb.java
deleted file mode 100644
index d6eccf94d3..0000000000
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/db/JobProfileDb.java
+++ /dev/null
@@ -1,253 +0,0 @@
-/*
- * 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.db;
-
-import org.apache.inlong.agent.conf.JobProfile;
-import org.apache.inlong.agent.constant.JobConstants;
-
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.ArrayList;
-import java.util.Arrays;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.Objects;
-import java.util.Set;
-import java.util.stream.Collectors;
-import java.util.stream.Stream;
-
-import static org.apache.inlong.agent.constant.JobConstants.JOB_ID;
-
-/**
- * Wrapper for job conf persistence.
- */
-public class JobProfileDb {
-
- private static final Logger LOGGER =
LoggerFactory.getLogger(JobProfileDb.class);
- private final Db db;
-
- public JobProfileDb(Db db) {
- this.db = db;
- }
-
- /**
- * get all restart jobs from db
- *
- * @Deprecated Use {@link JobProfileDb#getJobsByState(Set)}
- */
- @Deprecated
- public List<JobProfile> getRestartJobs() {
- List<JobProfile> jobsByState = getJobsByState(StateSearchKey.ACCEPTED);
- jobsByState.addAll(getJobsByState(StateSearchKey.RUNNING));
- LOGGER.info("try to get restart jobs from db {}", jobsByState);
- return jobsByState;
- }
-
- /**
- * update job state and search it by key name
- *
- * @param jobInstanceId job key name
- * @param stateSearchKey job state
- */
- public void updateJobState(String jobInstanceId, StateSearchKey
stateSearchKey) {
- KeyValueEntity entity = db.get(jobInstanceId);
- if (entity != null) {
- entity.setStateSearchKey(stateSearchKey);
- db.put(entity);
- }
- }
-
- /**
- * store job profile
- *
- * @param jobProfile job profile
- */
- public void storeJobFirstTime(JobProfile jobProfile) {
- if (jobProfile.allRequiredKeyExist()) {
- String keyName = jobProfile.get(JobConstants.JOB_INSTANCE_ID);
- jobProfile.setLong(JobConstants.JOB_STORE_TIME,
System.currentTimeMillis());
- KeyValueEntity entity = new KeyValueEntity(keyName,
- jobProfile.toJsonStr(),
jobProfile.get(JobConstants.JOB_DIR_FILTER_PATTERNS, ""));
- entity.setStateSearchKey(StateSearchKey.ACCEPTED);
- LOGGER.info("store job {} to db", jobProfile.toJsonStr());
- db.put(entity);
- }
- }
-
- /**
- * update job profile
- *
- * @param jobProfile
- */
- public void updateJobProfile(JobProfile jobProfile) {
- String instanceId = jobProfile.getInstanceId();
- KeyValueEntity entity = db.get(instanceId);
- if (entity == null) {
- LOGGER.warn("job profile {} doesn't exist, update job profile fail
{}", instanceId, jobProfile.toJsonStr());
- return;
- }
- entity.setJsonValue(jobProfile.toJsonStr());
- db.put(entity);
- }
-
- /**
- * check whether job is finished, note that non-exist job is regarded as
finished.
- */
- public boolean checkJobfinished(JobProfile jobProfile) {
- KeyValueEntity entity = db.get(jobProfile.getInstanceId());
- if (entity == null) {
- LOGGER.info("job profile {} doesn't exist",
jobProfile.getInstanceId());
- return true;
- }
- return entity.checkFinished();
- }
-
- /**
- * delete job by keyName
- */
- public void deleteJob(String keyName) {
- db.remove(keyName);
- }
-
- /**
- * get job profile by jobId
- *
- * @param jobId jobId
- * @return JobProfile
- */
- public JobProfile getJobById(String jobId) {
- KeyValueEntity keyValueEntity = db.get(jobId);
- if (keyValueEntity != null) {
- return keyValueEntity.getAsJobProfile();
- }
- return null;
- }
-
- /**
- * remove all expired jobs from db, including success and failed
- *
- * @param expireTime expireTime
- */
- public void removeExpireJobs(long expireTime) {
- // remove finished tasks
- List<KeyValueEntity> successEntityList =
db.search(StateSearchKey.SUCCESS);
- List<KeyValueEntity> failedEntityList =
db.search(StateSearchKey.FAILED);
- List<KeyValueEntity> entityList = new ArrayList<>(successEntityList);
- entityList.addAll(failedEntityList);
- for (KeyValueEntity entity : entityList) {
- if (entity.getKey().startsWith(JobConstants.JOB_ID_PREFIX)) {
- JobProfile profile = entity.getAsJobProfile();
- long storeTime = profile.getLong(JobConstants.JOB_STORE_TIME,
0);
- long currentTime = System.currentTimeMillis();
- if (storeTime == 0 || currentTime - storeTime > expireTime) {
- LOGGER.info("delete job {} because of timeout store time:
{}, expire time: {}",
- entity.getKey(), storeTime, expireTime);
- deleteJob(entity.getKey());
- }
- }
- }
- }
-
- /**
- * get job conf by state
- *
- * @param stateSearchKey state index for searching.
- * @return JobProfile
- */
- public JobProfile getJob(StateSearchKey stateSearchKey) {
- KeyValueEntity entity = db.searchOne(stateSearchKey);
- if (entity != null &&
entity.getKey().startsWith(JobConstants.JOB_ID_PREFIX)) {
- return entity.getAsJobProfile();
- }
- return null;
- }
-
- /**
- * get job reading specific file
- */
- public JobProfile getJobByFileName(String fileName) {
- KeyValueEntity entity = db.searchOne(fileName);
- if (entity != null &&
entity.getKey().startsWith(JobConstants.JOB_ID_PREFIX)) {
- return entity.getAsJobProfile();
- }
- return null;
- }
-
- /**
- * get list of job profiles.
- *
- * @param stateSearchKey state search key.
- * @return list of job profile.
- * @Deprecated Use {@link JobProfileDb#getJobsByState(Set)}
- */
- @Deprecated
- public List<JobProfile> getJobsByState(StateSearchKey stateSearchKey) {
- List<KeyValueEntity> entityList =
db.searchWithKeyPrefix(stateSearchKey, JobConstants.JOB_ID_PREFIX);
- List<JobProfile> profileList = new ArrayList<>();
- for (KeyValueEntity entity : entityList) {
- profileList.add(entity.getAsJobProfile());
- }
- return profileList;
- }
-
- /**
- * get list of job profiles by some state.
- *
- * @param stateSearchKeys state search keys.
- * @return list of job profile.
- */
- public List<JobProfile> getJobsByState(Set<StateSearchKey>
stateSearchKeys) {
- return stateSearchKeys.stream()
- .flatMap(stateSearchKey ->
db.searchWithKeyPrefix(stateSearchKey, JobConstants.JOB_ID_PREFIX).stream())
- .map(KeyValueEntity::getAsJobProfile)
- .collect(Collectors.toList());
- }
-
- /**
- * get all jobs.
- *
- * @return list of job profile.
- */
- public List<JobProfile> getAllJobs() {
- return
getJobsByState(Stream.of(StateSearchKey.values()).collect(Collectors.toSet()));
- }
-
- /**
- * check local job state from rocksDB.
- *
- * @return KV, key is job id and value is subtask config of job
- */
- public Map<String, List<String>> getJobsState() {
- List<KeyValueEntity> entityList =
db.search(Arrays.asList(StateSearchKey.values()));
- Map<String, List<String>> jobStateMap = new HashMap<>();
- for (KeyValueEntity entity : entityList) {
- List<String> tmpList = new ArrayList<>();
- JobProfile jobProfile = entity.getAsJobProfile();
- String jobState =
entity.getStateSearchKey().name().concat(":").concat(jobProfile.toJsonStr());
- tmpList.add(jobState);
- List<String> jobStates =
jobStateMap.putIfAbsent(jobProfile.get(JOB_ID), tmpList);
- if (Objects.nonNull(jobStates) && !jobStates.contains(jobState)) {
- jobStates.addAll(tmpList);
- jobStateMap.put(jobProfile.get(JOB_ID), jobStates);
- }
- }
- return jobStateMap;
- }
-}
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/BinlogJob.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/BinlogTask.java
similarity index 96%
rename from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/BinlogJob.java
rename to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/BinlogTask.java
index 86d41f2ad5..58a84be6e5 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/BinlogJob.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/BinlogTask.java
@@ -20,7 +20,7 @@ package org.apache.inlong.agent.pojo;
import lombok.Data;
@Data
-public class BinlogJob {
+public class BinlogTask {
private String user;
private String password;
@@ -60,7 +60,7 @@ public class BinlogJob {
}
@Data
- public static class BinlogJobTaskConfig {
+ public static class BinlogTaskConfig {
private String user;
private String password;
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/KafkaJob.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/KafkaTask.java
similarity index 96%
rename from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/KafkaJob.java
rename to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/KafkaTask.java
index c5ea65596e..963e89b4ab 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/KafkaJob.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/KafkaTask.java
@@ -20,7 +20,7 @@ package org.apache.inlong.agent.pojo;
import lombok.Data;
@Data
-public class KafkaJob {
+public class KafkaTask {
private String topic;
private String bootstrapServers;
@@ -63,7 +63,7 @@ public class KafkaJob {
}
@Data
- public static class KafkaJobTaskConfig {
+ public static class KafkaTaskConfig {
private String topic;
private String bootstrapServers;
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MongoJob.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MongoTask.java
similarity index 97%
rename from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MongoJob.java
rename to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MongoTask.java
index 1ecc4b6466..f0440feb8d 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MongoJob.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MongoTask.java
@@ -20,10 +20,10 @@ package org.apache.inlong.agent.pojo;
import lombok.Data;
/**
- * MongoJob : mongo job
+ * MongoTask : mongo task
*/
@Data
-public class MongoJob {
+public class MongoTask {
private String hosts;
private String user;
@@ -81,7 +81,7 @@ public class MongoJob {
}
@Data
- public static class MongoJobTaskConfig {
+ public static class MongoTaskConfig {
private String hosts;
private String username;
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MqttJob.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MqttTask.java
similarity index 95%
rename from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MqttJob.java
rename to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MqttTask.java
index 013b04a86e..10c1a03730 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MqttJob.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/MqttTask.java
@@ -20,10 +20,10 @@ package org.apache.inlong.agent.pojo;
import lombok.Data;
/**
- * MqttJob : mqtt job
+ * MqttTask : mqtt task
*/
@Data
-public class MqttJob {
+public class MqttTask {
private String serverURI;
private String userName;
@@ -39,7 +39,7 @@ public class MqttJob {
private String mqttVersion;
@Data
- public static class MqttJobConfig {
+ public static class MqttConfig {
private String serverURI;
private String username;
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/OracleJob.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/OracleTask.java
similarity index 90%
rename from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/OracleJob.java
rename to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/OracleTask.java
index ef2420b3a5..7810c0a807 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/OracleJob.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/OracleTask.java
@@ -20,7 +20,7 @@ package org.apache.inlong.agent.pojo;
import lombok.Data;
@Data
-public class OracleJob {
+public class OracleTask {
private String hostname;
private String user;
@@ -29,9 +29,9 @@ public class OracleJob {
private String serverName;
private String dbname;
- private OracleJob.Snapshot snapshot;
- private OracleJob.Offset offset;
- private OracleJob.History history;
+ private OracleTask.Snapshot snapshot;
+ private OracleTask.Offset offset;
+ private OracleTask.History history;
@Data
public static class Offset {
@@ -55,7 +55,7 @@ public class OracleJob {
}
@Data
- public static class OracleJobConfig {
+ public static class OracleTaskConfig {
private String hostname;
private String user;
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/PostgreSQLJob.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/PostgreSQLTask.java
similarity index 94%
rename from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/PostgreSQLJob.java
rename to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/PostgreSQLTask.java
index aea5e5fe04..0486eb4a7b 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/PostgreSQLJob.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/PostgreSQLTask.java
@@ -22,10 +22,10 @@ import lombok.Data;
import java.util.List;
/**
- * PostgreSQL Job info
+ * PostgreSQL Task info
*/
@Data
-public class PostgreSQLJob {
+public class PostgreSQLTask {
private String user;
private String password;
@@ -41,7 +41,7 @@ public class PostgreSQLJob {
private String primaryKey;
@Data
- public static class PostgreSQLJobConfig {
+ public static class PostgreSQLTaskConfig {
private String username;
private String password;
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/RedisJob.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/RedisTask.java
similarity index 95%
rename from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/RedisJob.java
rename to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/RedisTask.java
index 5917f9efe3..701bea2b8e 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/RedisJob.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/RedisTask.java
@@ -20,7 +20,7 @@ package org.apache.inlong.agent.pojo;
import lombok.Data;
@Data
-public class RedisJob {
+public class RedisTask {
private String authUser;
private String authPassword;
@@ -32,7 +32,7 @@ public class RedisJob {
private String replId;
@Data
- public static class RedisJobConfig {
+ public static class RedisTaskConfig {
private String username;
private String password;
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/SqlServerJob.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/SqlServerTask.java
similarity index 90%
rename from
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/SqlServerJob.java
rename to
inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/SqlServerTask.java
index c18ae32e57..56e4a9b920 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/SqlServerJob.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/SqlServerTask.java
@@ -20,7 +20,7 @@ package org.apache.inlong.agent.pojo;
import lombok.Data;
@Data
-public class SqlServerJob {
+public class SqlServerTask {
private String hostname;
private String user;
@@ -29,9 +29,9 @@ public class SqlServerJob {
private String serverName;
private String dbname;
- private SqlServerJob.Snapshot snapshot;
- private SqlServerJob.Offset offset;
- private SqlServerJob.History history;
+ private SqlServerTask.Snapshot snapshot;
+ private SqlServerTask.Offset offset;
+ private SqlServerTask.History history;
@Data
public static class Offset {
@@ -55,7 +55,7 @@ public class SqlServerJob {
}
@Data
- public static class SqlserverJobConfig {
+ public static class SqlserverTaskConfig {
private String hostname;
private String username;
diff --git
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/TaskProfileDto.java
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/TaskProfileDto.java
index d8bd595c6a..baf9e0c64b 100644
---
a/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/TaskProfileDto.java
+++
b/inlong-agent/agent-common/src/main/java/org/apache/inlong/agent/pojo/TaskProfileDto.java
@@ -20,9 +20,17 @@ package org.apache.inlong.agent.pojo;
import org.apache.inlong.agent.conf.AgentConfiguration;
import org.apache.inlong.agent.conf.TaskProfile;
import org.apache.inlong.agent.constant.CycleUnitType;
+import org.apache.inlong.agent.pojo.BinlogTask.BinlogTaskConfig;
import org.apache.inlong.agent.pojo.FileTask.FileTaskConfig;
import org.apache.inlong.agent.pojo.FileTask.Line;
+import org.apache.inlong.agent.pojo.KafkaTask.KafkaTaskConfig;
+import org.apache.inlong.agent.pojo.MongoTask.MongoTaskConfig;
+import org.apache.inlong.agent.pojo.MqttTask.MqttConfig;
+import org.apache.inlong.agent.pojo.OracleTask.OracleTaskConfig;
+import org.apache.inlong.agent.pojo.PostgreSQLTask.PostgreSQLTaskConfig;
import org.apache.inlong.agent.pojo.PulsarTask.PulsarTaskConfig;
+import org.apache.inlong.agent.pojo.RedisTask.RedisTaskConfig;
+import org.apache.inlong.agent.pojo.SqlServerTask.SqlserverTaskConfig;
import org.apache.inlong.common.constant.MQType;
import org.apache.inlong.common.enums.TaskTypeEnum;
import org.apache.inlong.common.pojo.agent.DataConfig;
@@ -91,44 +99,44 @@ public class TaskProfileDto {
private Task task;
private Proxy proxy;
- private static BinlogJob getBinlogJob(DataConfig dataConfigs) {
- BinlogJob.BinlogJobTaskConfig binlogJobTaskConfig =
GSON.fromJson(dataConfigs.getExtParams(),
- BinlogJob.BinlogJobTaskConfig.class);
+ private static BinlogTask getBinlogTask(DataConfig dataConfigs) {
+ BinlogTaskConfig binlogTaskConfig =
GSON.fromJson(dataConfigs.getExtParams(),
+ BinlogTaskConfig.class);
- BinlogJob binlogJob = new BinlogJob();
- binlogJob.setHostname(binlogJobTaskConfig.getHostname());
- binlogJob.setPassword(binlogJobTaskConfig.getPassword());
- binlogJob.setUser(binlogJobTaskConfig.getUser());
- binlogJob.setTableWhiteList(binlogJobTaskConfig.getTableWhiteList());
-
binlogJob.setDatabaseWhiteList(binlogJobTaskConfig.getDatabaseWhiteList());
- binlogJob.setSchema(binlogJobTaskConfig.getIncludeSchema());
- binlogJob.setPort(binlogJobTaskConfig.getPort());
- binlogJob.setOffsets(dataConfigs.getSnapshot());
- binlogJob.setDdl(binlogJobTaskConfig.getMonitoredDdl());
- binlogJob.setServerTimezone(binlogJobTaskConfig.getServerTimezone());
+ BinlogTask binlogTask = new BinlogTask();
+ binlogTask.setHostname(binlogTaskConfig.getHostname());
+ binlogTask.setPassword(binlogTaskConfig.getPassword());
+ binlogTask.setUser(binlogTaskConfig.getUser());
+ binlogTask.setTableWhiteList(binlogTaskConfig.getTableWhiteList());
+
binlogTask.setDatabaseWhiteList(binlogTaskConfig.getDatabaseWhiteList());
+ binlogTask.setSchema(binlogTaskConfig.getIncludeSchema());
+ binlogTask.setPort(binlogTaskConfig.getPort());
+ binlogTask.setOffsets(dataConfigs.getSnapshot());
+ binlogTask.setDdl(binlogTaskConfig.getMonitoredDdl());
+ binlogTask.setServerTimezone(binlogTaskConfig.getServerTimezone());
- BinlogJob.Offset offset = new BinlogJob.Offset();
- offset.setIntervalMs(binlogJobTaskConfig.getIntervalMs());
- offset.setFilename(binlogJobTaskConfig.getOffsetFilename());
-
offset.setSpecificOffsetFile(binlogJobTaskConfig.getSpecificOffsetFile());
-
offset.setSpecificOffsetPos(binlogJobTaskConfig.getSpecificOffsetPos());
+ BinlogTask.Offset offset = new BinlogTask.Offset();
+ offset.setIntervalMs(binlogTaskConfig.getIntervalMs());
+ offset.setFilename(binlogTaskConfig.getOffsetFilename());
+ offset.setSpecificOffsetFile(binlogTaskConfig.getSpecificOffsetFile());
+ offset.setSpecificOffsetPos(binlogTaskConfig.getSpecificOffsetPos());
- binlogJob.setOffset(offset);
+ binlogTask.setOffset(offset);
- BinlogJob.Snapshot snapshot = new BinlogJob.Snapshot();
- snapshot.setMode(binlogJobTaskConfig.getSnapshotMode());
+ BinlogTask.Snapshot snapshot = new BinlogTask.Snapshot();
+ snapshot.setMode(binlogTaskConfig.getSnapshotMode());
- binlogJob.setSnapshot(snapshot);
+ binlogTask.setSnapshot(snapshot);
- BinlogJob.History history = new BinlogJob.History();
- history.setFilename(binlogJobTaskConfig.getHistoryFilename());
+ BinlogTask.History history = new BinlogTask.History();
+ history.setFilename(binlogTaskConfig.getHistoryFilename());
- binlogJob.setHistory(history);
+ binlogTask.setHistory(history);
- return binlogJob;
+ return binlogTask;
}
- private static FileTask getFileJob(DataConfig dataConfig) {
+ private static FileTask getFileTask(DataConfig dataConfig) {
FileTask fileTask = new FileTask();
fileTask.setId(dataConfig.getTaskId());
@@ -185,32 +193,32 @@ public class TaskProfileDto {
return fileTask;
}
- private static KafkaJob getKafkaJob(DataConfig dataConfigs) {
-
- KafkaJob.KafkaJobTaskConfig kafkaJobTaskConfig =
GSON.fromJson(dataConfigs.getExtParams(),
- KafkaJob.KafkaJobTaskConfig.class);
- KafkaJob kafkaJob = new KafkaJob();
-
- KafkaJob.Bootstrap bootstrap = new KafkaJob.Bootstrap();
- bootstrap.setServers(kafkaJobTaskConfig.getBootstrapServers());
- kafkaJob.setBootstrap(bootstrap);
- KafkaJob.Partition partition = new KafkaJob.Partition();
- partition.setOffset(kafkaJobTaskConfig.getPartitionOffsets());
- kafkaJob.setPartition(partition);
- KafkaJob.Group group = new KafkaJob.Group();
- group.setId(kafkaJobTaskConfig.getGroupId());
- kafkaJob.setGroup(group);
- KafkaJob.RecordSpeed recordSpeed = new KafkaJob.RecordSpeed();
- recordSpeed.setLimit(kafkaJobTaskConfig.getRecordSpeedLimit());
- kafkaJob.setRecordSpeed(recordSpeed);
- KafkaJob.ByteSpeed byteSpeed = new KafkaJob.ByteSpeed();
- byteSpeed.setLimit(kafkaJobTaskConfig.getByteSpeedLimit());
- kafkaJob.setByteSpeed(byteSpeed);
- kafkaJob.setAutoOffsetReset(kafkaJobTaskConfig.getAutoOffsetReset());
-
- kafkaJob.setTopic(kafkaJobTaskConfig.getTopic());
-
- return kafkaJob;
+ private static KafkaTask getKafkaTask(DataConfig dataConfigs) {
+
+ KafkaTaskConfig kafkaTaskConfig =
GSON.fromJson(dataConfigs.getExtParams(),
+ KafkaTaskConfig.class);
+ KafkaTask kafkaTask = new KafkaTask();
+
+ KafkaTask.Bootstrap bootstrap = new KafkaTask.Bootstrap();
+ bootstrap.setServers(kafkaTaskConfig.getBootstrapServers());
+ kafkaTask.setBootstrap(bootstrap);
+ KafkaTask.Partition partition = new KafkaTask.Partition();
+ partition.setOffset(kafkaTaskConfig.getPartitionOffsets());
+ kafkaTask.setPartition(partition);
+ KafkaTask.Group group = new KafkaTask.Group();
+ group.setId(kafkaTaskConfig.getGroupId());
+ kafkaTask.setGroup(group);
+ KafkaTask.RecordSpeed recordSpeed = new KafkaTask.RecordSpeed();
+ recordSpeed.setLimit(kafkaTaskConfig.getRecordSpeedLimit());
+ kafkaTask.setRecordSpeed(recordSpeed);
+ KafkaTask.ByteSpeed byteSpeed = new KafkaTask.ByteSpeed();
+ byteSpeed.setLimit(kafkaTaskConfig.getByteSpeedLimit());
+ kafkaTask.setByteSpeed(byteSpeed);
+ kafkaTask.setAutoOffsetReset(kafkaTaskConfig.getAutoOffsetReset());
+
+ kafkaTask.setTopic(kafkaTaskConfig.getTopic());
+
+ return kafkaTask;
}
private static PulsarTask getPulsarTask(DataConfig dataConfig) {
@@ -230,163 +238,163 @@ public class TaskProfileDto {
return pulsarTask;
}
- private static PostgreSQLJob getPostgresJob(DataConfig dataConfigs) {
- PostgreSQLJob.PostgreSQLJobConfig config =
GSON.fromJson(dataConfigs.getExtParams(),
- PostgreSQLJob.PostgreSQLJobConfig.class);
- PostgreSQLJob postgreSQLJob = new PostgreSQLJob();
-
- postgreSQLJob.setUser(config.getUsername());
- postgreSQLJob.setPassword(config.getPassword());
- postgreSQLJob.setHostname(config.getHostname());
- postgreSQLJob.setPort(config.getPort());
- postgreSQLJob.setDbname(config.getDatabase());
- postgreSQLJob.setServername(config.getSchema());
- postgreSQLJob.setPluginname(config.getDecodingPluginName());
- postgreSQLJob.setTableNameList(config.getTableNameList());
- postgreSQLJob.setServerTimeZone(config.getServerTimeZone());
- postgreSQLJob.setScanStartupMode(config.getScanStartupMode());
- postgreSQLJob.setPrimaryKey(config.getPrimaryKey());
-
- return postgreSQLJob;
+ private static PostgreSQLTask getPostgresTask(DataConfig dataConfigs) {
+ PostgreSQLTaskConfig config = GSON.fromJson(dataConfigs.getExtParams(),
+ PostgreSQLTaskConfig.class);
+ PostgreSQLTask postgreSQLTask = new PostgreSQLTask();
+
+ postgreSQLTask.setUser(config.getUsername());
+ postgreSQLTask.setPassword(config.getPassword());
+ postgreSQLTask.setHostname(config.getHostname());
+ postgreSQLTask.setPort(config.getPort());
+ postgreSQLTask.setDbname(config.getDatabase());
+ postgreSQLTask.setServername(config.getSchema());
+ postgreSQLTask.setPluginname(config.getDecodingPluginName());
+ postgreSQLTask.setTableNameList(config.getTableNameList());
+ postgreSQLTask.setServerTimeZone(config.getServerTimeZone());
+ postgreSQLTask.setScanStartupMode(config.getScanStartupMode());
+ postgreSQLTask.setPrimaryKey(config.getPrimaryKey());
+
+ return postgreSQLTask;
}
- private static RedisJob getRedisJob(DataConfig dataConfig) {
- RedisJob.RedisJobConfig config =
GSON.fromJson(dataConfig.getExtParams(), RedisJob.RedisJobConfig.class);
- RedisJob redisJob = new RedisJob();
+ private static RedisTask getRedisTask(DataConfig dataConfig) {
+ RedisTaskConfig config = GSON.fromJson(dataConfig.getExtParams(),
RedisTaskConfig.class);
+ RedisTask redisTask = new RedisTask();
- redisJob.setAuthUser(config.getUsername());
- redisJob.setAuthPassword(config.getPassword());
- redisJob.setHostname(config.getHostname());
- redisJob.setPort(config.getPort());
- redisJob.setSsl(config.getSsl());
- redisJob.setReadTimeout(config.getTimeout());
- redisJob.setQueueSize(config.getQueueSize());
- redisJob.setReplId(config.getReplId());
+ redisTask.setAuthUser(config.getUsername());
+ redisTask.setAuthPassword(config.getPassword());
+ redisTask.setHostname(config.getHostname());
+ redisTask.setPort(config.getPort());
+ redisTask.setSsl(config.getSsl());
+ redisTask.setReadTimeout(config.getTimeout());
+ redisTask.setQueueSize(config.getQueueSize());
+ redisTask.setReplId(config.getReplId());
- return redisJob;
+ return redisTask;
}
- private static MongoJob getMongoJob(DataConfig dataConfigs) {
-
- MongoJob.MongoJobTaskConfig config =
GSON.fromJson(dataConfigs.getExtParams(),
- MongoJob.MongoJobTaskConfig.class);
- MongoJob mongoJob = new MongoJob();
-
- mongoJob.setHosts(config.getHosts());
- mongoJob.setUser(config.getUsername());
- mongoJob.setPassword(config.getPassword());
- mongoJob.setDatabaseIncludeList(config.getDatabaseIncludeList());
- mongoJob.setDatabaseExcludeList(config.getDatabaseExcludeList());
- mongoJob.setCollectionIncludeList(config.getCollectionIncludeList());
- mongoJob.setCollectionExcludeList(config.getCollectionExcludeList());
- mongoJob.setFieldExcludeList(config.getFieldExcludeList());
- mongoJob.setConnectTimeoutInMs(config.getConnectTimeoutInMs());
- mongoJob.setQueueSize(config.getQueueSize());
- mongoJob.setCursorMaxAwaitTimeInMs(config.getCursorMaxAwaitTimeInMs());
- mongoJob.setSocketTimeoutInMs(config.getSocketTimeoutInMs());
- mongoJob.setSelectionTimeoutInMs(config.getSelectionTimeoutInMs());
- mongoJob.setFieldRenames(config.getFieldRenames());
- mongoJob.setMembersAutoDiscover(config.getMembersAutoDiscover());
- mongoJob.setConnectMaxAttempts(config.getConnectMaxAttempts());
-
mongoJob.setConnectBackoffMaxDelayInMs(config.getConnectBackoffMaxDelayInMs());
-
mongoJob.setConnectBackoffInitialDelayInMs(config.getConnectBackoffInitialDelayInMs());
- mongoJob.setInitialSyncMaxThreads(config.getInitialSyncMaxThreads());
-
mongoJob.setSslInvalidHostnameAllowed(config.getSslInvalidHostnameAllowed());
- mongoJob.setSslEnabled(config.getSslEnabled());
- mongoJob.setPollIntervalInMs(config.getPollIntervalInMs());
-
- MongoJob.Offset offset = new MongoJob.Offset();
+ private static MongoTask getMongoTask(DataConfig dataConfigs) {
+
+ MongoTaskConfig config = GSON.fromJson(dataConfigs.getExtParams(),
+ MongoTaskConfig.class);
+ MongoTask mongoTask = new MongoTask();
+
+ mongoTask.setHosts(config.getHosts());
+ mongoTask.setUser(config.getUsername());
+ mongoTask.setPassword(config.getPassword());
+ mongoTask.setDatabaseIncludeList(config.getDatabaseIncludeList());
+ mongoTask.setDatabaseExcludeList(config.getDatabaseExcludeList());
+ mongoTask.setCollectionIncludeList(config.getCollectionIncludeList());
+ mongoTask.setCollectionExcludeList(config.getCollectionExcludeList());
+ mongoTask.setFieldExcludeList(config.getFieldExcludeList());
+ mongoTask.setConnectTimeoutInMs(config.getConnectTimeoutInMs());
+ mongoTask.setQueueSize(config.getQueueSize());
+
mongoTask.setCursorMaxAwaitTimeInMs(config.getCursorMaxAwaitTimeInMs());
+ mongoTask.setSocketTimeoutInMs(config.getSocketTimeoutInMs());
+ mongoTask.setSelectionTimeoutInMs(config.getSelectionTimeoutInMs());
+ mongoTask.setFieldRenames(config.getFieldRenames());
+ mongoTask.setMembersAutoDiscover(config.getMembersAutoDiscover());
+ mongoTask.setConnectMaxAttempts(config.getConnectMaxAttempts());
+
mongoTask.setConnectBackoffMaxDelayInMs(config.getConnectBackoffMaxDelayInMs());
+
mongoTask.setConnectBackoffInitialDelayInMs(config.getConnectBackoffInitialDelayInMs());
+ mongoTask.setInitialSyncMaxThreads(config.getInitialSyncMaxThreads());
+
mongoTask.setSslInvalidHostnameAllowed(config.getSslInvalidHostnameAllowed());
+ mongoTask.setSslEnabled(config.getSslEnabled());
+ mongoTask.setPollIntervalInMs(config.getPollIntervalInMs());
+
+ MongoTask.Offset offset = new MongoTask.Offset();
offset.setFilename(config.getOffsetFilename());
offset.setSpecificOffsetFile(config.getSpecificOffsetFile());
offset.setSpecificOffsetPos(config.getSpecificOffsetPos());
- mongoJob.setOffset(offset);
+ mongoTask.setOffset(offset);
- MongoJob.Snapshot snapshot = new MongoJob.Snapshot();
+ MongoTask.Snapshot snapshot = new MongoTask.Snapshot();
snapshot.setMode(config.getSnapshotMode());
- mongoJob.setSnapshot(snapshot);
+ mongoTask.setSnapshot(snapshot);
- MongoJob.History history = new MongoJob.History();
+ MongoTask.History history = new MongoTask.History();
history.setFilename(config.getHistoryFilename());
- mongoJob.setHistory(history);
+ mongoTask.setHistory(history);
- return mongoJob;
+ return mongoTask;
}
- private static OracleJob getOracleJob(DataConfig dataConfigs) {
- OracleJob.OracleJobConfig config =
GSON.fromJson(dataConfigs.getExtParams(),
- OracleJob.OracleJobConfig.class);
- OracleJob oracleJob = new OracleJob();
- oracleJob.setUser(config.getUser());
- oracleJob.setHostname(config.getHostname());
- oracleJob.setPassword(config.getPassword());
- oracleJob.setPort(config.getPort());
- oracleJob.setServerName(config.getServerName());
- oracleJob.setDbname(config.getDbname());
-
- OracleJob.Offset offset = new OracleJob.Offset();
+ private static OracleTask getOracleTask(DataConfig dataConfigs) {
+ OracleTaskConfig config = GSON.fromJson(dataConfigs.getExtParams(),
+ OracleTaskConfig.class);
+ OracleTask oracleTask = new OracleTask();
+ oracleTask.setUser(config.getUser());
+ oracleTask.setHostname(config.getHostname());
+ oracleTask.setPassword(config.getPassword());
+ oracleTask.setPort(config.getPort());
+ oracleTask.setServerName(config.getServerName());
+ oracleTask.setDbname(config.getDbname());
+
+ OracleTask.Offset offset = new OracleTask.Offset();
offset.setFilename(config.getOffsetFilename());
offset.setSpecificOffsetFile(config.getSpecificOffsetFile());
offset.setSpecificOffsetPos(config.getSpecificOffsetPos());
- oracleJob.setOffset(offset);
+ oracleTask.setOffset(offset);
- OracleJob.Snapshot snapshot = new OracleJob.Snapshot();
+ OracleTask.Snapshot snapshot = new OracleTask.Snapshot();
snapshot.setMode(config.getSnapshotMode());
- oracleJob.setSnapshot(snapshot);
+ oracleTask.setSnapshot(snapshot);
- OracleJob.History history = new OracleJob.History();
+ OracleTask.History history = new OracleTask.History();
history.setFilename(config.getHistoryFilename());
- oracleJob.setHistory(history);
+ oracleTask.setHistory(history);
- return oracleJob;
+ return oracleTask;
}
- private static SqlServerJob getSqlServerJob(DataConfig dataConfigs) {
- SqlServerJob.SqlserverJobConfig config =
GSON.fromJson(dataConfigs.getExtParams(),
- SqlServerJob.SqlserverJobConfig.class);
- SqlServerJob sqlServerJob = new SqlServerJob();
- sqlServerJob.setUser(config.getUsername());
- sqlServerJob.setHostname(config.getHostname());
- sqlServerJob.setPassword(config.getPassword());
- sqlServerJob.setPort(config.getPort());
- sqlServerJob.setServerName(config.getSchemaName());
- sqlServerJob.setDbname(config.getDatabase());
-
- SqlServerJob.Offset offset = new SqlServerJob.Offset();
+ private static SqlServerTask getSqlServerTask(DataConfig dataConfigs) {
+ SqlserverTaskConfig config = GSON.fromJson(dataConfigs.getExtParams(),
+ SqlserverTaskConfig.class);
+ SqlServerTask sqlServerTask = new SqlServerTask();
+ sqlServerTask.setUser(config.getUsername());
+ sqlServerTask.setHostname(config.getHostname());
+ sqlServerTask.setPassword(config.getPassword());
+ sqlServerTask.setPort(config.getPort());
+ sqlServerTask.setServerName(config.getSchemaName());
+ sqlServerTask.setDbname(config.getDatabase());
+
+ SqlServerTask.Offset offset = new SqlServerTask.Offset();
offset.setFilename(config.getOffsetFilename());
offset.setSpecificOffsetFile(config.getSpecificOffsetFile());
offset.setSpecificOffsetPos(config.getSpecificOffsetPos());
- sqlServerJob.setOffset(offset);
+ sqlServerTask.setOffset(offset);
- SqlServerJob.Snapshot snapshot = new SqlServerJob.Snapshot();
+ SqlServerTask.Snapshot snapshot = new SqlServerTask.Snapshot();
snapshot.setMode(config.getSnapshotMode());
- sqlServerJob.setSnapshot(snapshot);
+ sqlServerTask.setSnapshot(snapshot);
- SqlServerJob.History history = new SqlServerJob.History();
+ SqlServerTask.History history = new SqlServerTask.History();
history.setFilename(config.getHistoryFilename());
- sqlServerJob.setHistory(history);
+ sqlServerTask.setHistory(history);
- return sqlServerJob;
+ return sqlServerTask;
}
- public static MqttJob getMqttJob(DataConfig dataConfigs) {
- MqttJob.MqttJobConfig config =
GSON.fromJson(dataConfigs.getExtParams(),
- MqttJob.MqttJobConfig.class);
- MqttJob mqttJob = new MqttJob();
-
- mqttJob.setServerURI(config.getServerURI());
- mqttJob.setUserName(config.getUsername());
- mqttJob.setPassword(config.getPassword());
- mqttJob.setTopic(config.getTopic());
- mqttJob.setConnectionTimeOut(config.getConnectionTimeOut());
- mqttJob.setKeepAliveInterval(config.getKeepAliveInterval());
- mqttJob.setQos(config.getQos());
- mqttJob.setCleanSession(config.getCleanSession());
- mqttJob.setClientIdPrefix(config.getClientId());
- mqttJob.setQueueSize(config.getQueueSize());
- mqttJob.setAutomaticReconnect(config.getAutomaticReconnect());
- mqttJob.setMqttVersion(config.getMqttVersion());
-
- return mqttJob;
+ public static MqttTask getMqttTask(DataConfig dataConfigs) {
+ MqttConfig config = GSON.fromJson(dataConfigs.getExtParams(),
+ MqttConfig.class);
+ MqttTask mqttTask = new MqttTask();
+
+ mqttTask.setServerURI(config.getServerURI());
+ mqttTask.setUserName(config.getUsername());
+ mqttTask.setPassword(config.getPassword());
+ mqttTask.setTopic(config.getTopic());
+ mqttTask.setConnectionTimeOut(config.getConnectionTimeOut());
+ mqttTask.setKeepAliveInterval(config.getKeepAliveInterval());
+ mqttTask.setQos(config.getQos());
+ mqttTask.setCleanSession(config.getCleanSession());
+ mqttTask.setClientIdPrefix(config.getClientId());
+ mqttTask.setQueueSize(config.getQueueSize());
+ mqttTask.setAutomaticReconnect(config.getAutomaticReconnect());
+ mqttTask.setMqttVersion(config.getMqttVersion());
+
+ return mqttTask;
}
private static Proxy getProxy(DataConfig dataConfigs) {
@@ -457,14 +465,14 @@ public class TaskProfileDto {
switch (requireNonNull(taskType)) {
case SQL:
case BINLOG:
- BinlogJob binlogJob = getBinlogJob(dataConfig);
- task.setBinlogJob(binlogJob);
+ BinlogTask binlogTask = getBinlogTask(dataConfig);
+ task.setBinlogTask(binlogTask);
task.setSource(BINLOG_SOURCE);
profileDto.setTask(task);
break;
case FILE:
task.setTaskClass(DEFAULT_FILE_TASK);
- FileTask fileTask = getFileJob(dataConfig);
+ FileTask fileTask = getFileTask(dataConfig);
task.setCycleUnit(fileTask.getCycleUnit());
task.setFileTask(fileTask);
task.setSource(DEFAULT_SOURCE);
@@ -472,8 +480,8 @@ public class TaskProfileDto {
break;
case KAFKA:
task.setTaskClass(DEFAULT_KAFKA_TASK);
- KafkaJob kafkaJob = getKafkaJob(dataConfig);
- task.setKafkaJob(kafkaJob);
+ KafkaTask kafkaTask = getKafkaTask(dataConfig);
+ task.setKafkaTask(kafkaTask);
task.setSource(KAFKA_SOURCE);
profileDto.setTask(task);
break;
@@ -485,38 +493,38 @@ public class TaskProfileDto {
profileDto.setTask(task);
break;
case POSTGRES:
- PostgreSQLJob postgreSQLJob = getPostgresJob(dataConfig);
- task.setPostgreSQLJob(postgreSQLJob);
+ PostgreSQLTask postgreSQLTask = getPostgresTask(dataConfig);
+ task.setPostgreSQLTask(postgreSQLTask);
task.setSource(POSTGRESQL_SOURCE);
profileDto.setTask(task);
break;
case ORACLE:
- OracleJob oracleJob = getOracleJob(dataConfig);
- task.setOracleJob(oracleJob);
+ OracleTask oracleTask = getOracleTask(dataConfig);
+ task.setOracleTask(oracleTask);
task.setSource(ORACLE_SOURCE);
profileDto.setTask(task);
break;
case SQLSERVER:
- SqlServerJob sqlserverJob = getSqlServerJob(dataConfig);
- task.setSqlserverJob(sqlserverJob);
+ SqlServerTask sqlserverTask = getSqlServerTask(dataConfig);
+ task.setSqlserverTask(sqlserverTask);
task.setSource(SQLSERVER_SOURCE);
profileDto.setTask(task);
break;
case MONGODB:
- MongoJob mongoJob = getMongoJob(dataConfig);
- task.setMongoJob(mongoJob);
+ MongoTask mongoTask = getMongoTask(dataConfig);
+ task.setMongoTask(mongoTask);
task.setSource(MONGO_SOURCE);
profileDto.setTask(task);
break;
case REDIS:
- RedisJob redisJob = getRedisJob(dataConfig);
- task.setRedisJob(redisJob);
+ RedisTask redisTask = getRedisTask(dataConfig);
+ task.setRedisTask(redisTask);
task.setSource(REDIS_SOURCE);
profileDto.setTask(task);
break;
case MQTT:
- MqttJob mqttJob = getMqttJob(dataConfig);
- task.setMqttJob(mqttJob);
+ MqttTask mqttTask = getMqttTask(dataConfig);
+ task.setMqttTask(mqttTask);
task.setSource(MQTT_SOURCE);
profileDto.setTask(task);
break;
@@ -554,15 +562,15 @@ public class TaskProfileDto {
private String timeZone;
private FileTask fileTask;
- private BinlogJob binlogJob;
- private KafkaJob kafkaJob;
+ private BinlogTask binlogTask;
+ private KafkaTask kafkaTask;
private PulsarTask pulsarTask;
- private PostgreSQLJob postgreSQLJob;
- private OracleJob oracleJob;
- private MongoJob mongoJob;
- private RedisJob redisJob;
- private MqttJob mqttJob;
- private SqlServerJob sqlserverJob;
+ private PostgreSQLTask postgreSQLTask;
+ private OracleTask oracleTask;
+ private MongoTask mongoTask;
+ private RedisTask redisTask;
+ private MqttTask mqttTask;
+ private SqlServerTask sqlserverTask;
}
@Data
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sinks/filecollect/SenderManager.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sinks/filecollect/SenderManager.java
index 51056a1975..4e5c07ca26 100755
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sinks/filecollect/SenderManager.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sinks/filecollect/SenderManager.java
@@ -55,8 +55,8 @@ import static
org.apache.inlong.agent.constant.CommonConstants.PROXY_BATCH_FLUSH
import static
org.apache.inlong.agent.constant.FetcherConstants.AGENT_MANAGER_ADDR;
import static
org.apache.inlong.agent.constant.FetcherConstants.AGENT_MANAGER_AUTH_SECRET_ID;
import static
org.apache.inlong.agent.constant.FetcherConstants.AGENT_MANAGER_AUTH_SECRET_KEY;
-import static
org.apache.inlong.agent.constant.TaskConstants.DEFAULT_JOB_PROXY_SEND;
-import static org.apache.inlong.agent.constant.TaskConstants.JOB_PROXY_SEND;
+import static
org.apache.inlong.agent.constant.TaskConstants.DEFAULT_TASK_PROXY_SEND;
+import static org.apache.inlong.agent.constant.TaskConstants.TASK_PROXY_SEND;
import static
org.apache.inlong.agent.metrics.AgentMetricItem.KEY_INLONG_GROUP_ID;
import static
org.apache.inlong.agent.metrics.AgentMetricItem.KEY_INLONG_STREAM_ID;
import static org.apache.inlong.agent.metrics.AgentMetricItem.KEY_PLUGIN_ID;
@@ -110,7 +110,7 @@ public class SenderManager {
public SenderManager(InstanceProfile profile, String inlongGroupId, String
sourcePath) {
this.profile = profile;
managerAddr = agentConf.get(AGENT_MANAGER_ADDR);
- proxySend = profile.getBoolean(JOB_PROXY_SEND, DEFAULT_JOB_PROXY_SEND);
+ proxySend = profile.getBoolean(TASK_PROXY_SEND,
DEFAULT_TASK_PROXY_SEND);
totalAsyncBufSize = profile
.getInt(
CommonConstants.PROXY_TOTAL_ASYNC_PROXY_SIZE,
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/KafkaSource.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/KafkaSource.java
index 9acebc9ef7..fa9034c6bc 100644
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/KafkaSource.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/KafkaSource.java
@@ -65,14 +65,14 @@ import static
org.apache.inlong.agent.constant.CommonConstants.PROXY_PACKAGE_MAX
import static
org.apache.inlong.agent.constant.CommonConstants.PROXY_SEND_PARTITION_KEY;
import static
org.apache.inlong.agent.constant.FetcherConstants.AGENT_GLOBAL_READER_QUEUE_PERMIT;
import static
org.apache.inlong.agent.constant.FetcherConstants.AGENT_GLOBAL_READER_SOURCE_PERMIT;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_KAFKA_PARTITION_OFFSET_DELIMITER;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_OFFSET_DELIMITER;
import static org.apache.inlong.agent.constant.TaskConstants.OFFSET;
import static org.apache.inlong.agent.constant.TaskConstants.RESTORE_FROM_DB;
import static org.apache.inlong.agent.constant.TaskConstants.TASK_CYCLE_UNIT;
import static
org.apache.inlong.agent.constant.TaskConstants.TASK_KAFKA_AUTO_COMMIT_OFFSET_RESET;
import static
org.apache.inlong.agent.constant.TaskConstants.TASK_KAFKA_BOOTSTRAP_SERVERS;
import static org.apache.inlong.agent.constant.TaskConstants.TASK_KAFKA_OFFSET;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_KAFKA_OFFSET_DELIMITER;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_KAFKA_PARTITION_OFFSET_DELIMITER;
/**
* kafka source, split kafka source job into multi readers
@@ -149,10 +149,10 @@ public class KafkaSource extends AbstractSource {
isRestoreFromDB = profile.getBoolean(RESTORE_FROM_DB, false);
if (!isRestoreFromDB &&
StringUtils.isNotBlank(allPartitionOffsets)) {
// example:0#110_1#666_2#222
- String[] offsets =
allPartitionOffsets.split(JOB_OFFSET_DELIMITER);
+ String[] offsets =
allPartitionOffsets.split(TASK_KAFKA_OFFSET_DELIMITER);
for (String offset : offsets) {
-
partitionOffsets.put(Integer.valueOf(offset.split(JOB_KAFKA_PARTITION_OFFSET_DELIMITER)[0]),
-
Long.valueOf(offset.split(JOB_KAFKA_PARTITION_OFFSET_DELIMITER)[1]));
+
partitionOffsets.put(Integer.valueOf(offset.split(TASK_KAFKA_PARTITION_OFFSET_DELIMITER)[0]),
+
Long.valueOf(offset.split(TASK_KAFKA_PARTITION_OFFSET_DELIMITER)[1]));
}
}
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/LogFileSource.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/LogFileSource.java
index b6dd32b44d..3ba0f28f9c 100755
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/LogFileSource.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/LogFileSource.java
@@ -84,9 +84,9 @@ import static
org.apache.inlong.agent.constant.MetadataConstants.METADATA_FILE_N
import static
org.apache.inlong.agent.constant.MetadataConstants.METADATA_HOST_NAME;
import static
org.apache.inlong.agent.constant.MetadataConstants.METADATA_SOURCE_IP;
import static
org.apache.inlong.agent.constant.TaskConstants.DEFAULT_FILE_SOURCE_EXTEND_CLASS;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_FILE_META_ENV_LIST;
import static org.apache.inlong.agent.constant.TaskConstants.OFFSET;
import static org.apache.inlong.agent.constant.TaskConstants.TASK_CYCLE_UNIT;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_FILE_META_ENV_LIST;
/**
* Read text files
@@ -230,10 +230,10 @@ public class LogFileSource extends AbstractSource {
}
public void registerMeta(InstanceProfile jobConf) {
- if (!jobConf.hasKey(JOB_FILE_META_ENV_LIST)) {
+ if (!jobConf.hasKey(TASK_FILE_META_ENV_LIST)) {
return;
}
- String[] env = jobConf.get(JOB_FILE_META_ENV_LIST).split(COMMA);
+ String[] env = jobConf.get(TASK_FILE_META_ENV_LIST).split(COMMA);
Arrays.stream(env).forEach(data -> {
if (data.equalsIgnoreCase(KUBERNETES)) {
needMetadata = true;
@@ -248,8 +248,8 @@ public class LogFileSource extends AbstractSource {
}
private boolean isIncrement(InstanceProfile profile) {
- if (profile.hasKey(TaskConstants.JOB_FILE_CONTENT_COLLECT_TYPE) &&
DataCollectType.INCREMENT
-
.equalsIgnoreCase(profile.get(TaskConstants.JOB_FILE_CONTENT_COLLECT_TYPE))) {
+ if (profile.hasKey(TaskConstants.TASK_FILE_CONTENT_COLLECT_TYPE) &&
DataCollectType.INCREMENT
+
.equalsIgnoreCase(profile.get(TaskConstants.TASK_FILE_CONTENT_COLLECT_TYPE))) {
return true;
}
return false;
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/reader/MongoDBReader.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/reader/MongoDBReader.java
index f9036037a4..0c412cd4fe 100644
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/reader/MongoDBReader.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/sources/reader/MongoDBReader.java
@@ -78,34 +78,34 @@ import static
io.debezium.connector.mongodb.MongoDbConnectorConfig.SSL_ENABLED;
import static io.debezium.connector.mongodb.MongoDbConnectorConfig.USER;
import static
org.apache.inlong.agent.constant.CommonConstants.DEFAULT_MAP_CAPACITY;
import static org.apache.inlong.agent.constant.CommonConstants.PROXY_KEY_DATA;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_BACKOFF_INITIAL_DELAY;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_BACKOFF_MAX_DELAY;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_CAPTURE_MODE;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_COLLECTION_EXCLUDE_LIST;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_COLLECTION_INCLUDE_LIST;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_CONNECT_MAX_ATTEMPTS;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_CONNECT_TIMEOUT_MS;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_CURSOR_MAX_AWAIT;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_DATABASE_EXCLUDE_LIST;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_DATABASE_INCLUDE_LIST;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_FIELD_EXCLUDE_LIST;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_FIELD_RENAMES;
-import static org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_HOSTS;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_INITIAL_SYNC_MAX_THREADS;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_MEMBERS_DISCOVER;
-import static org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_OFFSETS;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_OFFSET_SPECIFIC_OFFSET_FILE;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_OFFSET_SPECIFIC_OFFSET_POS;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_PASSWORD;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_POLL_INTERVAL;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_QUEUE_SIZE;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_SELECTION_TIMEOUT;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_SNAPSHOT_MODE;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_SOCKET_TIMEOUT;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_SSL_ENABLE;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_SSL_INVALID_HOSTNAME_ALLOWED;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_STORE_HISTORY_FILENAME;
-import static org.apache.inlong.agent.constant.TaskConstants.JOB_MONGO_USER;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_BACKOFF_INITIAL_DELAY;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_BACKOFF_MAX_DELAY;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_CAPTURE_MODE;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_COLLECTION_EXCLUDE_LIST;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_COLLECTION_INCLUDE_LIST;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_CONNECT_MAX_ATTEMPTS;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_CONNECT_TIMEOUT_MS;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_CURSOR_MAX_AWAIT;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_DATABASE_EXCLUDE_LIST;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_DATABASE_INCLUDE_LIST;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_FIELD_EXCLUDE_LIST;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_FIELD_RENAMES;
+import static org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_HOSTS;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_INITIAL_SYNC_MAX_THREADS;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_MEMBERS_DISCOVER;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_OFFSETS;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_OFFSET_SPECIFIC_OFFSET_FILE;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_OFFSET_SPECIFIC_OFFSET_POS;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_PASSWORD;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_POLL_INTERVAL;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_QUEUE_SIZE;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_SELECTION_TIMEOUT;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_SNAPSHOT_MODE;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_SOCKET_TIMEOUT;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_SSL_ENABLE;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_SSL_INVALID_HOSTNAME_ALLOWED;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_STORE_HISTORY_FILENAME;
+import static org.apache.inlong.agent.constant.TaskConstants.TASK_MONGO_USER;
/**
* MongoDBReader : mongo source, split mongo source job into multi readers
@@ -223,18 +223,18 @@ public class MongoDBReader extends AbstractReader {
* @param jobConf job conf
*/
private void setGlobalParamsValue(InstanceProfile jobConf) {
- bufferPool = new
LinkedBlockingQueue<>(jobConf.getInt(JOB_MONGO_QUEUE_SIZE, 1000));
+ bufferPool = new
LinkedBlockingQueue<>(jobConf.getInt(TASK_MONGO_QUEUE_SIZE, 1000));
instanceId = jobConf.getInstanceId();
// offset file absolute path
- offsetStoreFileName = jobConf.get(JOB_MONGO_STORE_HISTORY_FILENAME,
+ offsetStoreFileName = jobConf.get(TASK_MONGO_STORE_HISTORY_FILENAME,
MongoDBSnapshotBase.getSnapshotFilePath()) + "/mongo-" +
instanceId + "-offset.dat";
// snapshot info
snapshot = new MongoDBSnapshotBase(offsetStoreFileName);
- String offset = jobConf.get(JOB_MONGO_OFFSETS, "");
+ String offset = jobConf.get(TASK_MONGO_OFFSETS, "");
snapshot.save(offset, new File(offsetStoreFileName));
// offset info
- specificOffsetFile =
jobConf.get(JOB_MONGO_OFFSET_SPECIFIC_OFFSET_FILE, "");
- specificOffsetPos = jobConf.get(JOB_MONGO_OFFSET_SPECIFIC_OFFSET_POS,
"-1");
+ specificOffsetFile =
jobConf.get(TASK_MONGO_OFFSET_SPECIFIC_OFFSET_FILE, "");
+ specificOffsetPos = jobConf.get(TASK_MONGO_OFFSET_SPECIFIC_OFFSET_POS,
"-1");
}
/**
@@ -275,31 +275,31 @@ public class MongoDBReader extends AbstractReader {
*/
private Properties buildMongoConnectorConfig(InstanceProfile jobConf) {
Configuration.Builder builder = Configuration.create();
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_HOSTS, HOSTS);
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_USER, USER);
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_PASSWORD,
PASSWORD);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_DATABASE_INCLUDE_LIST, DATABASE_INCLUDE_LIST);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_DATABASE_EXCLUDE_LIST, DATABASE_EXCLUDE_LIST);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_COLLECTION_INCLUDE_LIST, COLLECTION_INCLUDE_LIST);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_COLLECTION_EXCLUDE_LIST, COLLECTION_EXCLUDE_LIST);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_FIELD_EXCLUDE_LIST, FIELD_EXCLUDE_LIST);
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_SNAPSHOT_MODE,
SNAPSHOT_MODE);
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_CAPTURE_MODE,
CAPTURE_MODE);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_CONNECT_TIMEOUT_MS, CONNECT_TIMEOUT_MS);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_CURSOR_MAX_AWAIT, CURSOR_MAX_AWAIT_TIME_MS);
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_SOCKET_TIMEOUT,
SOCKET_TIMEOUT_MS);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_SELECTION_TIMEOUT, SERVER_SELECTION_TIMEOUT_MS);
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_FIELD_RENAMES,
FIELD_RENAMES);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_MEMBERS_DISCOVER, AUTO_DISCOVER_MEMBERS);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_CONNECT_MAX_ATTEMPTS, MAX_FAILED_CONNECTIONS);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_BACKOFF_MAX_DELAY, CONNECT_BACKOFF_MAX_DELAY_MS);
+ setEngineConfigIfNecessary(jobConf, builder, TASK_MONGO_HOSTS, HOSTS);
+ setEngineConfigIfNecessary(jobConf, builder, TASK_MONGO_USER, USER);
+ setEngineConfigIfNecessary(jobConf, builder, TASK_MONGO_PASSWORD,
PASSWORD);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_DATABASE_INCLUDE_LIST, DATABASE_INCLUDE_LIST);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_DATABASE_EXCLUDE_LIST, DATABASE_EXCLUDE_LIST);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_COLLECTION_INCLUDE_LIST, COLLECTION_INCLUDE_LIST);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_COLLECTION_EXCLUDE_LIST, COLLECTION_EXCLUDE_LIST);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_FIELD_EXCLUDE_LIST, FIELD_EXCLUDE_LIST);
+ setEngineConfigIfNecessary(jobConf, builder, TASK_MONGO_SNAPSHOT_MODE,
SNAPSHOT_MODE);
+ setEngineConfigIfNecessary(jobConf, builder, TASK_MONGO_CAPTURE_MODE,
CAPTURE_MODE);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_CONNECT_TIMEOUT_MS, CONNECT_TIMEOUT_MS);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_CURSOR_MAX_AWAIT, CURSOR_MAX_AWAIT_TIME_MS);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_SOCKET_TIMEOUT, SOCKET_TIMEOUT_MS);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_SELECTION_TIMEOUT, SERVER_SELECTION_TIMEOUT_MS);
+ setEngineConfigIfNecessary(jobConf, builder, TASK_MONGO_FIELD_RENAMES,
FIELD_RENAMES);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_MEMBERS_DISCOVER, AUTO_DISCOVER_MEMBERS);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_CONNECT_MAX_ATTEMPTS, MAX_FAILED_CONNECTIONS);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_BACKOFF_MAX_DELAY, CONNECT_BACKOFF_MAX_DELAY_MS);
setEngineConfigIfNecessary(jobConf, builder,
- JOB_MONGO_BACKOFF_INITIAL_DELAY,
CONNECT_BACKOFF_INITIAL_DELAY_MS);
- setEngineConfigIfNecessary(jobConf, builder,
JOB_MONGO_INITIAL_SYNC_MAX_THREADS, MAX_COPY_THREADS);
+ TASK_MONGO_BACKOFF_INITIAL_DELAY,
CONNECT_BACKOFF_INITIAL_DELAY_MS);
+ setEngineConfigIfNecessary(jobConf, builder,
TASK_MONGO_INITIAL_SYNC_MAX_THREADS, MAX_COPY_THREADS);
setEngineConfigIfNecessary(jobConf, builder,
- JOB_MONGO_SSL_INVALID_HOSTNAME_ALLOWED,
SSL_ALLOW_INVALID_HOSTNAMES);
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_SSL_ENABLE,
SSL_ENABLED);
- setEngineConfigIfNecessary(jobConf, builder, JOB_MONGO_POLL_INTERVAL,
MONGODB_POLL_INTERVAL_MS);
+ TASK_MONGO_SSL_INVALID_HOSTNAME_ALLOWED,
SSL_ALLOW_INVALID_HOSTNAMES);
+ setEngineConfigIfNecessary(jobConf, builder, TASK_MONGO_SSL_ENABLE,
SSL_ENABLED);
+ setEngineConfigIfNecessary(jobConf, builder, TASK_MONGO_POLL_INTERVAL,
MONGODB_POLL_INTERVAL_MS);
Properties props = builder.build().asProperties();
props.setProperty("offset.storage.file.filename", offsetStoreFileName);
@@ -307,12 +307,12 @@ public class MongoDBReader extends AbstractReader {
props.setProperty("name", "engine-" + instanceId);
props.setProperty("mongodb.name", "inlong-mongodb-" + instanceId);
- String snapshotMode = props.getOrDefault(JOB_MONGO_SNAPSHOT_MODE,
"").toString();
+ String snapshotMode = props.getOrDefault(TASK_MONGO_SNAPSHOT_MODE,
"").toString();
if (Objects.equals(SnapshotModeConstants.INITIAL, snapshotMode)) {
- Preconditions.checkNotNull(JOB_MONGO_OFFSET_SPECIFIC_OFFSET_FILE,
- JOB_MONGO_OFFSET_SPECIFIC_OFFSET_FILE + " cannot be null");
- Preconditions.checkNotNull(JOB_MONGO_OFFSET_SPECIFIC_OFFSET_POS,
- JOB_MONGO_OFFSET_SPECIFIC_OFFSET_POS + " cannot be null");
+ Preconditions.checkNotNull(TASK_MONGO_OFFSET_SPECIFIC_OFFSET_FILE,
+ TASK_MONGO_OFFSET_SPECIFIC_OFFSET_FILE + " cannot be
null");
+ Preconditions.checkNotNull(TASK_MONGO_OFFSET_SPECIFIC_OFFSET_POS,
+ TASK_MONGO_OFFSET_SPECIFIC_OFFSET_POS + " cannot be null");
props.setProperty("offset.storage",
InLongFileOffsetBackingStore.class.getCanonicalName());
props.setProperty(InLongFileOffsetBackingStore.OFFSET_STATE_VALUE,
serializeOffset(instanceId, specificOffsetFile,
specificOffsetPos));
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/MetaDataUtils.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/MetaDataUtils.java
index aa215bb539..863eaf7952 100644
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/MetaDataUtils.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/MetaDataUtils.java
@@ -37,8 +37,8 @@ import static
org.apache.inlong.agent.constant.KubernetesConstants.CONTAINER_ID;
import static
org.apache.inlong.agent.constant.KubernetesConstants.CONTAINER_NAME;
import static org.apache.inlong.agent.constant.KubernetesConstants.NAMESPACE;
import static org.apache.inlong.agent.constant.KubernetesConstants.POD_NAME;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_FILE_META_FILTER_BY_LABELS;
-import static
org.apache.inlong.agent.constant.TaskConstants.JOB_FILE_PROPERTIES;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_FILE_META_FILTER_BY_LABELS;
+import static
org.apache.inlong.agent.constant.TaskConstants.TASK_FILE_PROPERTIES;
/**
* Metadata utils
@@ -90,20 +90,20 @@ public class MetaDataUtils {
* get labels of pod
*/
public static Map<String, String> getPodLabels(AbstractConfiguration
taskProfile) {
- if (Objects.isNull(taskProfile) ||
!taskProfile.hasKey(JOB_FILE_META_FILTER_BY_LABELS)) {
+ if (Objects.isNull(taskProfile) ||
!taskProfile.hasKey(TASK_FILE_META_FILTER_BY_LABELS)) {
return new HashMap<>();
}
- String labels = taskProfile.get(JOB_FILE_META_FILTER_BY_LABELS);
+ String labels = taskProfile.get(TASK_FILE_META_FILTER_BY_LABELS);
Type type = new TypeToken<HashMap<String, String>>() {
}.getType();
return GSON.fromJson(labels, type);
}
public static List<String> getNamespace(AbstractConfiguration taskProfile)
{
- if (Objects.isNull(taskProfile) ||
!taskProfile.hasKey(JOB_FILE_PROPERTIES)) {
+ if (Objects.isNull(taskProfile) ||
!taskProfile.hasKey(TASK_FILE_PROPERTIES)) {
return null;
}
- String property = taskProfile.get(JOB_FILE_PROPERTIES);
+ String property = taskProfile.get(TASK_FILE_PROPERTIES);
Type type = new TypeToken<HashMap<Integer, String>>() {
}.getType();
Map<String, String> properties = GSON.fromJson(property, type);
@@ -121,10 +121,10 @@ public class MetaDataUtils {
* get name of pod
*/
public static String getPodName(AbstractConfiguration taskProfile) {
- if (Objects.isNull(taskProfile) ||
!taskProfile.hasKey(JOB_FILE_PROPERTIES)) {
+ if (Objects.isNull(taskProfile) ||
!taskProfile.hasKey(TASK_FILE_PROPERTIES)) {
return null;
}
- String property = taskProfile.get(JOB_FILE_PROPERTIES);
+ String property = taskProfile.get(TASK_FILE_PROPERTIES);
Type type = new TypeToken<HashMap<Integer, String>>() {
}.getType();
Map<String, String> properties = GSON.fromJson(property, type);
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/PluginUtils.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/PluginUtils.java
index e08313c10b..db0c1ce6d3 100755
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/PluginUtils.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/PluginUtils.java
@@ -55,8 +55,8 @@ import static
org.apache.inlong.agent.constant.KubernetesConstants.HTTPS;
import static
org.apache.inlong.agent.constant.KubernetesConstants.KUBERNETES_SERVICE_HOST;
import static
org.apache.inlong.agent.constant.KubernetesConstants.KUBERNETES_SERVICE_PORT;
import static
org.apache.inlong.agent.constant.TaskConstants.FILE_DIR_FILTER_PATTERNS;
-import static org.apache.inlong.agent.constant.TaskConstants.JOB_RETRY_TIME;
import static
org.apache.inlong.agent.constant.TaskConstants.TASK_FILE_TIME_OFFSET;
+import static org.apache.inlong.agent.constant.TaskConstants.TASK_RETRY_TIME;
/**
* Utils for plugin package.
@@ -107,10 +107,10 @@ public class PluginUtils {
* if the job is retry job, the date is determined
*/
public static void updateRetryTime(TaskProfile jobConf,
Collection<PathPattern> patterns) {
- if (jobConf.hasKey(JOB_RETRY_TIME)) {
+ if (jobConf.hasKey(TASK_RETRY_TIME)) {
LOGGER.info("job {} is retry job with specific time, update file
time to {}"
- + "", jobConf.toJsonStr(), jobConf.get(JOB_RETRY_TIME));
- patterns.forEach(pattern ->
pattern.updateDateFormatRegex(jobConf.get(JOB_RETRY_TIME)));
+ + "", jobConf.toJsonStr(), jobConf.get(TASK_RETRY_TIME));
+ patterns.forEach(pattern ->
pattern.updateDateFormatRegex(jobConf.get(TASK_RETRY_TIME)));
}
}
@@ -121,7 +121,7 @@ public class PluginUtils {
TaskProfile copiedProfile =
TaskProfile.parseJsonStr(taskProfile.toJsonStr());
String md5 = AgentUtils.getFileMd5(pendingFile);
copiedProfile.set(pendingFile.getAbsolutePath() + ".md5", md5);
- copiedProfile.set(TaskConstants.JOB_FILE_TRIGGER, null); // del
trigger id
+ copiedProfile.set(TaskConstants.TASK_FILE_TRIGGER, null); // del
trigger id
copiedProfile.set(TaskConstants.FILE_DIR_FILTER_PATTERNS,
pendingFile.getAbsolutePath());
return copiedProfile;
}
diff --git
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/RocksDBUtils.java
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/RocksDBUtils.java
index 278b0c9f22..14a86dba65 100644
---
a/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/RocksDBUtils.java
+++
b/inlong-agent/agent-plugins/src/main/java/org/apache/inlong/agent/plugin/utils/RocksDBUtils.java
@@ -24,12 +24,9 @@ import org.apache.inlong.agent.constant.TaskConstants;
import org.apache.inlong.agent.db.Db;
import org.apache.inlong.agent.db.RocksDbImp;
import org.apache.inlong.agent.db.TaskProfileDb;
-import org.apache.inlong.agent.utils.AgentUtils;
import java.util.List;
-import static org.apache.inlong.agent.constant.TaskConstants.JOB_ID_PREFIX;
-
public class RocksDBUtils {
public static void main(String[] args) {
@@ -49,9 +46,6 @@ public class RocksDBUtils {
triggerProfile.set(TaskConstants.TASK_DIR_FILTER_PATTERN,
null);
}
- triggerProfile.set(TaskConstants.JOB_INSTANCE_ID,
- AgentUtils.getSingleJobId(JOB_ID_PREFIX,
triggerProfile.getTaskId()));
-
triggerProfileDb.storeTask(triggerProfile);
});
}
diff --git
a/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/sinks/MockSink.java
b/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/sinks/MockSink.java
index 0ff07d0e42..86431c4ee7 100644
---
a/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/sinks/MockSink.java
+++
b/inlong-agent/agent-plugins/src/test/java/org/apache/inlong/agent/plugin/sinks/MockSink.java
@@ -29,7 +29,7 @@ import java.util.List;
import java.util.concurrent.atomic.AtomicLong;
import static org.apache.inlong.agent.constant.JobConstants.JOB_CYCLE_UNIT;
-import static org.apache.inlong.agent.constant.JobConstants.JOB_DATA_TIME;
+import static org.apache.inlong.agent.constant.TaskConstants.SINK_DATA_TIME;
public class MockSink extends AbstractSink {
@@ -58,7 +58,7 @@ public class MockSink extends AbstractSink {
@Override
public void init(InstanceProfile jobConf) {
super.init(jobConf);
- dataTime =
AgentUtils.timeStrConvertToMillSec(jobConf.get(JOB_DATA_TIME, ""),
+ dataTime =
AgentUtils.timeStrConvertToMillSec(jobConf.get(SINK_DATA_TIME, ""),
jobConf.get(JOB_CYCLE_UNIT, ""));
sourceFileName = "test";
LOGGER.info("get dataTime is : {}", dataTime);