From 0ae5c4c31c86d645a323418e3d5849626e272159 Mon Sep 17 00:00:00 2001 From: "longpeng.zlp" Date: Thu, 24 Oct 2024 17:28:37 +0800 Subject: [PATCH 1/2] refactor base task, remove runtime info to TaskContext --- .../odc/service/task/TaskApplicationTest.java | 2 +- .../agent/runtime/ExecutorRequestHandler.java | 15 +-- .../odc/agent/runtime/TaskExecutor.java | 2 +- .../odc/agent/runtime}/TaskMonitor.java | 7 +- .../odc/agent/runtime/TaskRuntimeInfo.java | 33 +++++ .../agent/runtime/ThreadPoolTaskExecutor.java | 126 ++++++++++++++++-- .../objectstorage/ObjectStorageHandler.java | 29 ++-- .../util/ObjectStorageUtils.java | 4 +- .../{executor => }/ExceptionListener.java | 2 +- .../odc/service/task/SharedStorage.java | 77 +++++++++++ .../com/oceanbase/odc/service/task/Task.java | 1 - .../task/{executor => }/TaskContext.java | 16 ++- .../odc/service/task/TaskEventListener.java | 55 ++++++++ .../odc/service/task/base/BaseTask.java | 64 +++------ .../base/dataarchive/DataArchiveTask.java | 2 +- .../databasechange/DatabaseChangeTask.java | 16 +-- .../LogicalDatabaseChangeTask.java | 2 +- .../task/base/precheck/PreCheckTask.java | 4 +- .../task/base/rollback/RollbackPlanTask.java | 14 +- .../task/base/sqlplan/SqlPlanTask.java | 21 +-- .../odc/service/task/dummy/DummyTask.java | 2 +- .../task/executor/task/BaseTaskTest.java | 20 ++- 22 files changed, 389 insertions(+), 125 deletions(-) rename server/{odc-service/src/main/java/com/oceanbase/odc/service/task/executor => odc-server/src/main/java/com/oceanbase/odc/agent/runtime}/TaskMonitor.java (96%) create mode 100644 server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskRuntimeInfo.java rename server/odc-service/src/main/java/com/oceanbase/odc/service/task/{executor => }/ExceptionListener.java (93%) create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/SharedStorage.java rename server/odc-service/src/main/java/com/oceanbase/odc/service/task/{executor => }/TaskContext.java (77%) create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/TaskEventListener.java diff --git a/server/integration-test/src/test/java/com/oceanbase/odc/service/task/TaskApplicationTest.java b/server/integration-test/src/test/java/com/oceanbase/odc/service/task/TaskApplicationTest.java index 1b0ba39479..b649662c03 100644 --- a/server/integration-test/src/test/java/com/oceanbase/odc/service/task/TaskApplicationTest.java +++ b/server/integration-test/src/test/java/com/oceanbase/odc/service/task/TaskApplicationTest.java @@ -79,7 +79,7 @@ private void assertRunningResult(JobIdentity ji) { try { Thread.sleep(60 * 1000L); - Task task = ThreadPoolTaskExecutor.getInstance().getTask(ji); + Task task = ThreadPoolTaskExecutor.getInstance().getTaskRuntimeInfo(ji); Assert.assertSame(JobStatus.DONE, task.getStatus()); } catch (InterruptedException e) { throw new RuntimeException(e); diff --git a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ExecutorRequestHandler.java b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ExecutorRequestHandler.java index db3192e397..7dcc8cca35 100644 --- a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ExecutorRequestHandler.java +++ b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ExecutorRequestHandler.java @@ -25,12 +25,9 @@ import com.oceanbase.odc.service.common.response.Responses; import com.oceanbase.odc.service.common.response.SuccessResponse; import com.oceanbase.odc.service.common.util.UrlUtils; -import com.oceanbase.odc.service.task.Task; -import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.constants.JobExecutorUrls; import com.oceanbase.odc.service.task.executor.DefaultTaskResult; import com.oceanbase.odc.service.task.executor.DefaultTaskResultBuilder; -import com.oceanbase.odc.service.task.executor.TaskMonitor; import com.oceanbase.odc.service.task.executor.logger.LogBiz; import com.oceanbase.odc.service.task.executor.logger.LogBizImpl; import com.oceanbase.odc.service.task.executor.logger.LogUtils; @@ -89,21 +86,21 @@ public SuccessResponse process(HttpMethod httpMethod, String uri, String matcher = modifyParametersPattern.matcher(path); if (matcher.find()) { JobIdentity ji = getJobIdentity(matcher); - Task task = ThreadPoolTaskExecutor.getInstance().getTask(ji); - boolean result = task.modify(JobUtils.fromJsonToMap(requestData)); + TaskRuntimeInfo runtimeInfo = ThreadPoolTaskExecutor.getInstance().getTaskRuntimeInfo(ji); + boolean result = runtimeInfo.getTask().modify(JobUtils.fromJsonToMap(requestData)); return Responses.ok(result); } matcher = getResultPattern.matcher(path); if (matcher.find()) { JobIdentity ji = getJobIdentity(matcher); - BaseTask task = ThreadPoolTaskExecutor.getInstance().getTask(ji); - TaskMonitor taskMonitor = task.getTaskMonitor(); - DefaultTaskResult result = DefaultTaskResultBuilder.build(task); + TaskRuntimeInfo runtimeInfo = ThreadPoolTaskExecutor.getInstance().getTaskRuntimeInfo(ji); + TaskMonitor taskMonitor = runtimeInfo.getTaskMonitor(); + DefaultTaskResult result = DefaultTaskResultBuilder.build(runtimeInfo.getTask()); if (taskMonitor != null && MapUtils.isNotEmpty(taskMonitor.getLogMetadata())) { result.setLogMetadata(taskMonitor.getLogMetadata()); // assign final error message - DefaultTaskResultBuilder.assignErrorMessage(result, task); + DefaultTaskResultBuilder.assignErrorMessage(result, runtimeInfo.getTask()); taskMonitor.markLogMetaCollected(); } DefaultTaskResult copiedResult = ObjectUtil.deepCopy(result, DefaultTaskResult.class); diff --git a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskExecutor.java b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskExecutor.java index a749f42a62..8ab11d2a0f 100644 --- a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskExecutor.java +++ b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskExecutor.java @@ -29,6 +29,6 @@ public interface TaskExecutor { boolean cancel(JobIdentity ji); - BaseTask getTask(JobIdentity ji); + TaskRuntimeInfo getTaskRuntimeInfo(JobIdentity ji); } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/TaskMonitor.java b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskMonitor.java similarity index 96% rename from server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/TaskMonitor.java rename to server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskMonitor.java index bfcc4f4749..a93e3c8627 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/TaskMonitor.java +++ b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskMonitor.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package com.oceanbase.odc.service.task.executor; +package com.oceanbase.odc.agent.runtime; import java.util.HashMap; import java.util.Map; @@ -35,6 +35,11 @@ import com.oceanbase.odc.service.task.constants.JobParametersKeyConstants; import com.oceanbase.odc.service.task.constants.JobServerUrls; import com.oceanbase.odc.service.task.enums.JobStatus; +import com.oceanbase.odc.service.task.executor.DefaultTaskResult; +import com.oceanbase.odc.service.task.executor.DefaultTaskResultBuilder; +import com.oceanbase.odc.service.task.executor.HeartbeatRequest; +import com.oceanbase.odc.service.task.executor.TaskReporter; +import com.oceanbase.odc.service.task.executor.TraceDecoratorThreadFactory; import com.oceanbase.odc.service.task.executor.logger.LogBizImpl; import com.oceanbase.odc.service.task.util.JobUtils; diff --git a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskRuntimeInfo.java b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskRuntimeInfo.java new file mode 100644 index 0000000000..8cdc5f5997 --- /dev/null +++ b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/TaskRuntimeInfo.java @@ -0,0 +1,33 @@ +/* + * Copyright (c) 2023 OceanBase. + * + * Licensed 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 com.oceanbase.odc.agent.runtime; + +import java.util.concurrent.Future; + +import com.oceanbase.odc.service.task.base.BaseTask; + +import lombok.Data; + +/** + * @author longpeng.zlp + * @date 2024/10/24 14:32 + */ +@Data +public class TaskRuntimeInfo { + private final BaseTask task; + private final Future future; + private final TaskMonitor taskMonitor; +} diff --git a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ThreadPoolTaskExecutor.java b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ThreadPoolTaskExecutor.java index 87df7d85ae..479e0cafc6 100644 --- a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ThreadPoolTaskExecutor.java +++ b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ThreadPoolTaskExecutor.java @@ -15,8 +15,13 @@ */ package com.oceanbase.odc.agent.runtime; -import java.util.HashMap; +import java.io.File; +import java.io.IOException; +import java.io.InputStream; +import java.net.URL; import java.util.Map; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -27,13 +32,19 @@ import com.oceanbase.odc.common.concurrent.ExecutorUtils; import com.oceanbase.odc.core.shared.PreConditions; import com.oceanbase.odc.core.task.TaskThreadFactory; +import com.oceanbase.odc.service.objectstorage.cloud.CloudObjectStorageService; +import com.oceanbase.odc.service.objectstorage.cloud.model.ObjectStorageConfiguration; +import com.oceanbase.odc.service.task.ExceptionListener; +import com.oceanbase.odc.service.task.SharedStorage; import com.oceanbase.odc.service.task.Task; +import com.oceanbase.odc.service.task.TaskContext; +import com.oceanbase.odc.service.task.TaskEventListener; import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.executor.TraceDecoratorThreadFactory; -import com.oceanbase.odc.service.task.executor.task.ExceptionListener; -import com.oceanbase.odc.service.task.executor.task.TaskContext; import com.oceanbase.odc.service.task.schedule.JobIdentity; +import com.oceanbase.odc.service.task.util.CloudObjectStorageServiceBuilder; +import com.oceanbase.odc.service.task.util.JobUtils; import lombok.extern.slf4j.Slf4j; @@ -47,8 +58,7 @@ public class ThreadPoolTaskExecutor implements TaskExecutor { private static final TaskExecutor TASK_EXECUTOR = new ThreadPoolTaskExecutor(); - private final Map> tasks = new HashMap<>(); - private final Map> futures = new HashMap<>(); + private final Map tasks = new ConcurrentHashMap<>(); private final ExecutorService executor; private ThreadPoolTaskExecutor() { @@ -68,6 +78,9 @@ synchronized public void execute(BaseTask task, JobContext jc) { if (tasks.containsKey(jobIdentity)) { throw new IllegalArgumentException("Task already exists, jobIdentity=" + jobIdentity.getId()); } + // init cloud objet storage service and task monitor + CloudObjectStorageService cloudObjectStorageService = buildCloudStorageService(jc); + TaskMonitor taskMonitor = new TaskMonitor(task, cloudObjectStorageService); Future future = executor.submit(() -> { try { task.start(new TaskContext() { @@ -80,19 +93,108 @@ public ExceptionListener getExceptionListener() { public JobContext getJobContext() { return jc; } + + @Override + public TaskEventListener getTaskEventListener() { + return taskEventListener(taskMonitor); + } + + @Override + public SharedStorage getSharedStorage() { + return sharedStorage(cloudObjectStorageService); + } }); } catch (Exception e) { log.error("Task start failed, jobIdentity={}.", jobIdentity.getId(), e); task.onException(e); } }); - futures.put(jobIdentity, future); - tasks.put(jobIdentity, task); + tasks.put(jobIdentity, new TaskRuntimeInfo(task, future, taskMonitor)); + } + + private SharedStorage sharedStorage(CloudObjectStorageService cloudObjectStorageService) { + return new SharedStorage() { + @Override + public boolean available() { + return null != cloudObjectStorageService && cloudObjectStorageService.supported(); + } + + @Override + public InputStream download(String objectId) throws IOException { + return cloudObjectStorageService.getObject(objectId); + } + + @Override + public String upload(String expectName, File localFile) throws IOException { + return cloudObjectStorageService.upload(expectName, localFile); + } + + @Override + public String uploadTemp(String expectName, File localFile) throws IOException { + return cloudObjectStorageService.uploadTemp(expectName, localFile); + } + + @Override + public URL getDownloadURL(String objectId) throws IOException { + return cloudObjectStorageService.generateDownloadUrl(objectId); + } + + @Override + public String getPathPrefix() { + return cloudObjectStorageService.getBucketName(); + } + }; + } + + /** + * build task event listener + * + * @param taskMonitor + * @return + */ + private TaskEventListener taskEventListener(TaskMonitor taskMonitor) { + return new TaskEventListener() { + @Override + public void onTaskStart(Task task) { + taskMonitor.monitor(); + } + + @Override + public void onTaskStop(Task task) {} + + @Override + public void onTaskModify(Task task) {} + + @Override + public void onTaskFinalize(Task task) { + taskMonitor.finalWork(); + } + }; + } + + /** + * build task monitor + * + * @param jobContext + * @return + */ + protected CloudObjectStorageService buildCloudStorageService(JobContext jobContext) { + Optional storageConfig = JobUtils.getObjectStorageConfiguration(); + CloudObjectStorageService cloudObjectStorageService = null; + try { + if (storageConfig.isPresent()) { + cloudObjectStorageService = CloudObjectStorageServiceBuilder.build(storageConfig.get()); + } + } catch (Throwable e) { + log.warn("Init cloud object storage service failed, id={}.", jobContext.getJobIdentity().getId(), e); + } + return cloudObjectStorageService; } @Override public boolean cancel(JobIdentity ji) { - Task task = getTask(ji); + TaskRuntimeInfo runtimeInfo = getTaskRuntimeInfo(ji); + BaseTask task = runtimeInfo.getTask(); Future stopFuture = executor.submit(task::stop); boolean result = false; try { @@ -117,9 +219,9 @@ public boolean cancel(JobIdentity ji) { } @Override - public BaseTask getTask(JobIdentity ji) { - BaseTask task = tasks.get(ji); - PreConditions.notNull(task, "task", "Task not found, jobIdentity=" + ji.getId()); - return task; + public TaskRuntimeInfo getTaskRuntimeInfo(JobIdentity ji) { + TaskRuntimeInfo runtimeInfo = tasks.get(ji); + PreConditions.notNull(runtimeInfo, "task", "Task not found, jobIdentity=" + ji.getId()); + return runtimeInfo; } } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/ObjectStorageHandler.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/ObjectStorageHandler.java index cb92d36e66..d9dc071cd6 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/ObjectStorageHandler.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/ObjectStorageHandler.java @@ -15,21 +15,18 @@ */ package com.oceanbase.odc.service.objectstorage; -import java.io.File; -import java.io.FileInputStream; import java.io.IOException; -import java.nio.charset.StandardCharsets; +import java.io.InputStream; import org.apache.commons.io.FileUtils; -import org.apache.commons.io.IOUtils; import org.springframework.core.io.Resource; import com.oceanbase.odc.core.shared.exception.InternalServerError; -import com.oceanbase.odc.service.objectstorage.cloud.CloudObjectStorageService; import com.oceanbase.odc.service.objectstorage.model.ObjectMetadata; import com.oceanbase.odc.service.objectstorage.model.StorageObject; import com.oceanbase.odc.service.objectstorage.operator.LocalFileOperator; import com.oceanbase.odc.service.objectstorage.util.ObjectStorageUtils; +import com.oceanbase.odc.service.task.SharedStorage; import lombok.extern.slf4j.Slf4j; @@ -42,27 +39,22 @@ public class ObjectStorageHandler { private final LocalFileOperator localFileOperator; - private final CloudObjectStorageService cloudObjectStorageService; + private final SharedStorage sharedStorage; - public ObjectStorageHandler(CloudObjectStorageService cloudObjectStorageService, String localDir) { + public ObjectStorageHandler(SharedStorage sharedStorage, String localDir) { this.localFileOperator = new LocalFileOperator(localDir); - this.cloudObjectStorageService = cloudObjectStorageService; - } - - public String loadObjectContentAsString(ObjectMetadata metadata) throws IOException { - StorageObject storageObject = loadObject(metadata); - return IOUtils.toString(storageObject.getContent(), StandardCharsets.UTF_8); + this.sharedStorage = sharedStorage; } public StorageObject loadObject(ObjectMetadata metadata) { Resource resource; try { - if (cloudObjectStorageService.supported()) { + if (sharedStorage.available()) { // OSS supported, load file from local if exists, otherwise load file from oss if (localFileOperator.isLocalFileAbsent(metadata)) { log.info("File is absent in local, load file from oss, bucket={}, objectId={}", metadata.getBucketName(), metadata.getObjectId()); - loadObjectFromOss(metadata); + loadObjectFromSharedObject(metadata); } resource = localFileOperator.loadAsResource(metadata.getBucketName(), metadata.getObjectId()); } else { @@ -78,14 +70,11 @@ public StorageObject loadObject(ObjectMetadata metadata) { } } - private void loadObjectFromOss(ObjectMetadata metadata) throws IOException { - File tempFile = cloudObjectStorageService.downloadToTempFile(metadata.getObjectId()); - try (FileInputStream inputStream = new FileInputStream(tempFile)) { + private void loadObjectFromSharedObject(ObjectMetadata metadata) throws IOException { + try (InputStream inputStream = sharedStorage.download(metadata.getObjectId())) { FileUtils.copyInputStreamToFile(inputStream, localFileOperator.getOrCreateLocalFile(metadata.getBucketName(), metadata.getObjectId())); - } finally { - FileUtils.deleteQuietly(tempFile); } } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/util/ObjectStorageUtils.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/util/ObjectStorageUtils.java index e369e375d2..03b24cd2f0 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/util/ObjectStorageUtils.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/objectstorage/util/ObjectStorageUtils.java @@ -25,9 +25,9 @@ import com.oceanbase.odc.core.shared.Verify; import com.oceanbase.odc.service.flow.task.model.SizeAwareInputStream; import com.oceanbase.odc.service.objectstorage.ObjectStorageHandler; -import com.oceanbase.odc.service.objectstorage.cloud.CloudObjectStorageService; import com.oceanbase.odc.service.objectstorage.model.ObjectMetadata; import com.oceanbase.odc.service.objectstorage.model.StorageObject; +import com.oceanbase.odc.service.task.SharedStorage; import lombok.NonNull; import lombok.extern.slf4j.Slf4j; @@ -49,7 +49,7 @@ public static String concatObjectId(String objectId, String extension) { } public static SizeAwareInputStream loadObjectsForTask(@NonNull List metadatas, - CloudObjectStorageService cloudOSS, String executorDataPath, long maxReadBytes) throws IOException { + SharedStorage cloudOSS, String executorDataPath, long maxReadBytes) throws IOException { InputStream inputStream = new ByteArrayInputStream(new byte[0]); long totalBytes = 0; for (ObjectMetadata metadata : metadatas) { diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/ExceptionListener.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/ExceptionListener.java similarity index 93% rename from server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/ExceptionListener.java rename to server/odc-service/src/main/java/com/oceanbase/odc/service/task/ExceptionListener.java index 7aac18956a..6d2bed6f90 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/ExceptionListener.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/ExceptionListener.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package com.oceanbase.odc.service.task.executor.task; +package com.oceanbase.odc.service.task; /** * listener for exception diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/SharedStorage.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/SharedStorage.java new file mode 100644 index 0000000000..1006cfdf93 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/SharedStorage.java @@ -0,0 +1,77 @@ +/* + * Copyright (c) 2023 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task; + +import java.io.File; +import java.io.IOException; +import java.io.InputStream; +import java.net.URL; + +/** + * @author longpeng.zlp + * @date 2024/10/24 14:52 + */ +public interface SharedStorage { + + /** + * if the shared storage is available to user + * + * @return + */ + boolean available(); + + /** + * create input stream by storage object id + * + * @param objectId in shared storage + * @return + */ + InputStream download(String objectId) throws IOException; + + /** + * upload permanent file to shared storage + * + * @param expectName + * @param localFile + * @return object id for shared storage + */ + String upload(String expectName, File localFile) throws IOException; + + /** + * upload expired file to shared storage + * + * @param expectName + * @param localFile + * @return object id for shared storage + */ + String uploadTemp(String expectName, File localFile) throws IOException; + + + /** + * translate objectId to download url + * + * @param objectId + * @return + */ + URL getDownloadURL(String objectId) throws IOException; + + /** + * get prefix path on shared storage + * + * @return + */ + String getPathPrefix(); +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/Task.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/Task.java index 06d12b89e2..7dc72ea65f 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/Task.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/Task.java @@ -19,7 +19,6 @@ import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.enums.JobStatus; -import com.oceanbase.odc.service.task.executor.task.TaskContext; /** * Task interface. Each task should implement this interface diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/TaskContext.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/TaskContext.java similarity index 77% rename from server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/TaskContext.java rename to server/odc-service/src/main/java/com/oceanbase/odc/service/task/TaskContext.java index e5971bb421..1d2354a809 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/executor/TaskContext.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/TaskContext.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package com.oceanbase.odc.service.task.executor.task; +package com.oceanbase.odc.service.task; import com.oceanbase.odc.service.task.caller.JobContext; @@ -37,4 +37,18 @@ public interface TaskContext { * @return */ JobContext getJobContext(); + + /** + * provide task event listener + * + * @return + */ + TaskEventListener getTaskEventListener(); + + /** + * get shared storage for task upload or download file + * + * @return + */ + SharedStorage getSharedStorage(); } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/TaskEventListener.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/TaskEventListener.java new file mode 100644 index 0000000000..e68a2d8c97 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/TaskEventListener.java @@ -0,0 +1,55 @@ +/* + * Copyright (c) 2023 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task; + +/** + * listen task event + * + * @author longpeng.zlp + * @date 2024/10/24 11:51 + */ +public interface TaskEventListener { + /** + * call when task start called + * + * @param task + */ + void onTaskStart(Task task); + + + /** + * call when task stop called + * + * @param task + */ + void onTaskStop(Task task); + + + /** + * call when task modify called + * + * @param task + */ + void onTaskModify(Task task); + + /** + * call when task final closed + * + * @param task + */ + void onTaskFinalize(Task task); + +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/BaseTask.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/BaseTask.java index 1ac410c920..816f1434ae 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/BaseTask.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/BaseTask.java @@ -18,23 +18,16 @@ import java.util.Collections; import java.util.Map; import java.util.Objects; -import java.util.Optional; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; -import com.oceanbase.odc.service.objectstorage.cloud.CloudObjectStorageService; -import com.oceanbase.odc.service.objectstorage.cloud.model.ObjectStorageConfiguration; +import com.oceanbase.odc.service.task.ExceptionListener; import com.oceanbase.odc.service.task.Task; +import com.oceanbase.odc.service.task.TaskContext; import com.oceanbase.odc.service.task.caller.DefaultJobContext; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.enums.JobStatus; -import com.oceanbase.odc.service.task.executor.TaskMonitor; -import com.oceanbase.odc.service.task.executor.task.ExceptionListener; -import com.oceanbase.odc.service.task.executor.task.TaskContext; -import com.oceanbase.odc.service.task.util.CloudObjectStorageServiceBuilder; -import com.oceanbase.odc.service.task.util.JobUtils; -import lombok.Getter; import lombok.extern.slf4j.Slf4j; /** @@ -45,37 +38,27 @@ public abstract class BaseTask implements Task, ExceptionListener { private final AtomicBoolean closed = new AtomicBoolean(false); - private JobContext context; + protected TaskContext context; private Map jobParameters; private volatile JobStatus status = JobStatus.PREPARING; - private CloudObjectStorageService cloudObjectStorageService; // only save latest exception if any // it will be cleaned if been fetched protected AtomicReference latestException = new AtomicReference<>(); - @Getter - private TaskMonitor taskMonitor; - @Override public void start(TaskContext taskContext) { - this.context = taskContext.getJobContext(); - log.info("Start task, id={}.", context.getJobIdentity().getId()); + this.context = taskContext; + JobContext jobContext = taskContext.getJobContext(); + log.info("Start task, id={}.", jobContext.getJobIdentity().getId()); - this.jobParameters = Collections.unmodifiableMap(context.getJobParameters()); - log.info("Init task parameters success, id={}.", context.getJobIdentity().getId()); + this.jobParameters = Collections.unmodifiableMap(jobContext.getJobParameters()); + log.info("Init task parameters success, id={}.", jobContext.getJobIdentity().getId()); try { - initCloudObjectStorageService(); - } catch (Exception e) { - log.warn("Init cloud object storage service failed, id={}.", getJobId(), e); - } - - this.taskMonitor = createTaskMonitor(); - try { - doInit(context); + doInit(context.getJobContext()); updateStatus(JobStatus.RUNNING); - taskMonitor.monitor(); - if (doStart(context, taskContext)) { + context.getTaskEventListener().onTaskStart(this); + if (doStart(context.getJobContext(), taskContext)) { updateStatus(JobStatus.DONE); } else { updateStatus(JobStatus.FAILED); @@ -89,10 +72,6 @@ public void start(TaskContext taskContext) { } } - protected TaskMonitor createTaskMonitor() { - return new TaskMonitor(this, cloudObjectStorageService); - } - @Override public boolean stop() { try { @@ -103,6 +82,9 @@ public boolean stop() { // doRefresh cannot execute if update status to 'canceled'. updateStatus(JobStatus.CANCELING); } + if (null != context) { + context.getTaskEventListener().onTaskStop(this); + } return true; } catch (Throwable e) { log.warn("Stop task failed, id={}", getJobId(), e); @@ -130,14 +112,12 @@ public boolean modify(Map jobParameters) { } catch (Exception e) { log.warn("Do after modified job parameters failed", e); } + if (null != context) { + context.getTaskEventListener().onTaskModify(this); + } return true; } - private void initCloudObjectStorageService() { - Optional storageConfig = JobUtils.getObjectStorageConfiguration(); - storageConfig.ifPresent(osc -> this.cloudObjectStorageService = CloudObjectStorageServiceBuilder.build(osc)); - } - private void close() { if (closed.compareAndSet(false, true)) { try { @@ -146,14 +126,12 @@ private void close() { // do nothing } log.info("Task completed, id={}, status={}.", getJobId(), getStatus()); - taskMonitor.finalWork(); + if (null != context) { + context.getTaskEventListener().onTaskFinalize(this); + } } } - protected CloudObjectStorageService getCloudObjectStorageService() { - return cloudObjectStorageService; - } - @Override public JobStatus getStatus() { return status; @@ -161,7 +139,7 @@ public JobStatus getStatus() { @Override public JobContext getJobContext() { - return context; + return context.getJobContext(); } protected void updateStatus(JobStatus status) { diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/dataarchive/DataArchiveTask.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/dataarchive/DataArchiveTask.java index 7e1666f637..611fabe809 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/dataarchive/DataArchiveTask.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/dataarchive/DataArchiveTask.java @@ -35,10 +35,10 @@ import com.oceanbase.odc.service.dlm.utils.DlmJobIdUtil; import com.oceanbase.odc.service.schedule.job.DLMJobReq; import com.oceanbase.odc.service.schedule.model.DlmTableUnitStatistic; +import com.oceanbase.odc.service.task.TaskContext; import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.constants.JobParametersKeyConstants; -import com.oceanbase.odc.service.task.executor.task.TaskContext; import com.oceanbase.odc.service.task.util.JobUtils; import com.oceanbase.tools.migrator.common.enums.JobType; import com.oceanbase.tools.migrator.job.Job; diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/databasechange/DatabaseChangeTask.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/databasechange/DatabaseChangeTask.java index afe0d13967..f273cfb4fb 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/databasechange/DatabaseChangeTask.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/databasechange/DatabaseChangeTask.java @@ -85,18 +85,18 @@ import com.oceanbase.odc.service.flow.task.model.DatabaseChangeParameters; import com.oceanbase.odc.service.flow.task.model.DatabaseChangeResult; import com.oceanbase.odc.service.flow.task.model.SizeAwareInputStream; -import com.oceanbase.odc.service.objectstorage.cloud.CloudObjectStorageService; import com.oceanbase.odc.service.objectstorage.util.ObjectStorageUtils; import com.oceanbase.odc.service.session.OdcStatementCallBack; import com.oceanbase.odc.service.session.factory.DefaultConnectSessionFactory; import com.oceanbase.odc.service.session.initializer.ConsoleTimeoutInitializer; import com.oceanbase.odc.service.session.model.SqlExecuteResult; +import com.oceanbase.odc.service.task.SharedStorage; +import com.oceanbase.odc.service.task.TaskContext; import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.constants.JobParametersKeyConstants; import com.oceanbase.odc.service.task.constants.JobServerUrls; import com.oceanbase.odc.service.task.exception.JobException; -import com.oceanbase.odc.service.task.executor.task.TaskContext; import com.oceanbase.odc.service.task.util.HttpClientUtils; import com.oceanbase.odc.service.task.util.JobUtils; import com.oceanbase.tools.dbbrowser.parser.ParserUtil; @@ -141,7 +141,7 @@ public class DatabaseChangeTask extends BaseTask { private long taskId; @Override - protected void doInit(JobContext context) { + protected void doInit(JobContext jobContext) { taskId = getJobContext().getJobIdentity().getId(); log.info("Initiating database change task, taskId={}", taskId); this.parameters = JobUtils.fromJson(getJobParameters().get(JobParametersKeyConstants.TASK_PARAMETER_JSON_KEY), @@ -162,7 +162,7 @@ protected void doInit(JobContext context) { try { SizeAwareInputStream sizeAwareInputStream = ObjectStorageUtils.loadObjectsForTask(this.parameters.getSqlFileObjectMetadatas(), - getCloudObjectStorageService(), JobUtils.getExecutorDataPath(), -1); + context.getSharedStorage(), JobUtils.getExecutorDataPath(), -1); sqlTotalBytes += sizeAwareInputStream.getTotalBytes(); sqlInputStream = sizeAwareInputStream.getInputStream(); } catch (IOException exception) { @@ -429,12 +429,12 @@ private void writeZipFile() { OdcFileUtil.zip(String.format(zipFileRootPath), String.format("%s.zip", zipFileRootPath)); log.info("Database change task result set was saved as local zip file, file name={}", zipFileId); // Public cloud scenario, need to upload files to OSS - CloudObjectStorageService cloudObjectStorageService = getCloudObjectStorageService(); - if (Objects.nonNull(cloudObjectStorageService) && cloudObjectStorageService.supported()) { + SharedStorage sharedStorage = context.getSharedStorage(); + if (sharedStorage.available()) { File tempZipFile = new File(String.format("%s.zip", zipFileRootPath)); try { - String objectName = cloudObjectStorageService.uploadTemp(zipFileId + ".zip", tempZipFile); - zipFileDownloadUrl = cloudObjectStorageService.getBucketName() + "/" + objectName; + String objectName = sharedStorage.uploadTemp(zipFileId + ".zip", tempZipFile); + zipFileDownloadUrl = sharedStorage.getPathPrefix() + "/" + objectName; log.info("Upload database change task result set zip file to OSS, file name={}", zipFileId); } catch (Exception exception) { log.warn("Upload database change task result set zip file to OSS failed, file name={}", zipFileId); diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/logicdatabasechange/LogicalDatabaseChangeTask.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/logicdatabasechange/LogicalDatabaseChangeTask.java index aaccd731e9..c6961e9b8b 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/logicdatabasechange/LogicalDatabaseChangeTask.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/logicdatabasechange/LogicalDatabaseChangeTask.java @@ -52,10 +52,10 @@ import com.oceanbase.odc.service.connection.logicaldatabase.model.DetailLogicalTableResp; import com.oceanbase.odc.service.schedule.model.PublishLogicalDatabaseChangeReq; import com.oceanbase.odc.service.session.model.SqlExecuteResult; +import com.oceanbase.odc.service.task.TaskContext; import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.constants.JobParametersKeyConstants; -import com.oceanbase.odc.service.task.executor.task.TaskContext; import com.oceanbase.tools.dbbrowser.parser.SqlParser; import com.oceanbase.tools.loaddump.utils.CollectionUtils; import com.oceanbase.tools.sqlparser.statement.Statement; diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/precheck/PreCheckTask.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/precheck/PreCheckTask.java index c00ec4078f..cee3873bb0 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/precheck/PreCheckTask.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/precheck/PreCheckTask.java @@ -62,10 +62,10 @@ import com.oceanbase.odc.service.sqlcheck.SqlCheckRuleFactory; import com.oceanbase.odc.service.sqlcheck.model.CheckViolation; import com.oceanbase.odc.service.sqlcheck.rule.SqlCheckRules; +import com.oceanbase.odc.service.task.TaskContext; import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.constants.JobParametersKeyConstants; -import com.oceanbase.odc.service.task.executor.task.TaskContext; import com.oceanbase.odc.service.task.util.JobUtils; import lombok.NonNull; @@ -197,7 +197,7 @@ private void loadUploadFileInputStream() throws IOException { List objectMetadataList = this.parameters.getSqlFileObjectMetadatas(); if (Objects.nonNull(params) && CollectionUtils.isNotEmpty(objectMetadataList)) { this.uploadFileInputStream = ObjectStorageUtils.loadObjectsForTask(objectMetadataList, - getCloudObjectStorageService(), JobUtils.getExecutorDataPath(), -1).getInputStream(); + context.getSharedStorage(), JobUtils.getExecutorDataPath(), -1).getInputStream(); this.uploadFileSqlIterator = SqlUtils.iterator(this.parameters.getConnectionConfig().getDialectType(), params.getDelimiter(), this.uploadFileInputStream, StandardCharsets.UTF_8); } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/rollback/RollbackPlanTask.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/rollback/RollbackPlanTask.java index 843689aba1..dea874b9a1 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/rollback/RollbackPlanTask.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/rollback/RollbackPlanTask.java @@ -39,7 +39,6 @@ import com.oceanbase.odc.service.common.util.SqlUtils; import com.oceanbase.odc.service.connection.model.ConnectionConfig; import com.oceanbase.odc.service.flow.task.model.RollbackPlanTaskResult; -import com.oceanbase.odc.service.objectstorage.cloud.CloudObjectStorageService; import com.oceanbase.odc.service.objectstorage.model.ObjectMetadata; import com.oceanbase.odc.service.objectstorage.util.ObjectStorageUtils; import com.oceanbase.odc.service.rollbackplan.GenerateRollbackPlan; @@ -47,10 +46,11 @@ import com.oceanbase.odc.service.rollbackplan.UnsupportedSqlTypeForRollbackPlanException; import com.oceanbase.odc.service.rollbackplan.model.RollbackPlan; import com.oceanbase.odc.service.session.factory.DefaultConnectSessionFactory; +import com.oceanbase.odc.service.task.SharedStorage; +import com.oceanbase.odc.service.task.TaskContext; import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.constants.JobParametersKeyConstants; -import com.oceanbase.odc.service.task.executor.task.TaskContext; import com.oceanbase.odc.service.task.util.JobUtils; import lombok.extern.slf4j.Slf4j; @@ -200,7 +200,7 @@ private void loadUploadFileInputStream() throws IOException { if (CollectionUtils.isNotEmpty(objectMetadataList)) { this.uploadFileInputStream = ObjectStorageUtils - .loadObjectsForTask(objectMetadataList, getCloudObjectStorageService(), + .loadObjectsForTask(objectMetadataList, context.getSharedStorage(), JobUtils.getExecutorDataPath(), parameters.getRollbackProperties().getMaxRollbackContentSizeBytes()) .getInputStream(); @@ -223,12 +223,12 @@ private void handleRollbackResult(String rollbackResult) { String resultFileId = StringUtils.uuid(); String filePath = String.format("%s/%s.sql", resultFileRootPath, resultFileId); FileUtils.writeStringToFile(new File(filePath), rollbackResult, StandardCharsets.UTF_8); - CloudObjectStorageService cloudObjectStorageService = getCloudObjectStorageService(); - if (Objects.nonNull(cloudObjectStorageService) && cloudObjectStorageService.supported()) { + SharedStorage sharedStorage = context.getSharedStorage(); + if (sharedStorage.available()) { File tempFile = new File(filePath); try { - String objectName = cloudObjectStorageService.uploadTemp(resultFileId + ".sql", tempFile); - resultFileDownloadUrl = cloudObjectStorageService.getBucketName() + "/" + objectName; + String objectName = sharedStorage.uploadTemp(resultFileId + ".sql", tempFile); + resultFileDownloadUrl = sharedStorage.getPathPrefix() + "/" + objectName; log.info("Upload generated rollback plan task result file to OSS, file name={}", resultFileId); } finally { OdcFileUtil.deleteFiles(tempFile); diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/sqlplan/SqlPlanTask.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/sqlplan/SqlPlanTask.java index 13be0b19ea..f059664d67 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/sqlplan/SqlPlanTask.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/base/sqlplan/SqlPlanTask.java @@ -60,17 +60,17 @@ import com.oceanbase.odc.service.common.util.OdcFileUtil; import com.oceanbase.odc.service.common.util.SqlUtils; import com.oceanbase.odc.service.connection.model.ConnectionConfig; -import com.oceanbase.odc.service.objectstorage.cloud.CloudObjectStorageService; import com.oceanbase.odc.service.schedule.job.PublishSqlPlanJobReq; import com.oceanbase.odc.service.session.OdcStatementCallBack; import com.oceanbase.odc.service.session.factory.DefaultConnectSessionFactory; import com.oceanbase.odc.service.session.initializer.ConsoleTimeoutInitializer; import com.oceanbase.odc.service.session.model.SqlExecuteResult; import com.oceanbase.odc.service.sqlplan.model.SqlPlanTaskResult; +import com.oceanbase.odc.service.task.SharedStorage; +import com.oceanbase.odc.service.task.TaskContext; import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.constants.JobParametersKeyConstants; -import com.oceanbase.odc.service.task.executor.task.TaskContext; import com.oceanbase.odc.service.task.util.JobUtils; import com.oceanbase.tools.dbbrowser.parser.ParserUtil; import com.oceanbase.tools.dbbrowser.parser.constant.GeneralSqlType; @@ -189,14 +189,14 @@ private void initSqlIterator() { return; } - CloudObjectStorageService cloudObjectStorageService = getCloudObjectStorageService(); - if (Objects.isNull(cloudObjectStorageService) || !cloudObjectStorageService.supported()) { + SharedStorage sharedStorage = context.getSharedStorage(); + if (!sharedStorage.available()) { log.warn("Cloud object storage service not supported."); throw new UnexpectedException("Cloud object storage service not supported"); } for (String sqlObjectId : parameters.getSqlObjectIds()) { - try (InputStream current = cloudObjectStorageService.getObject(sqlObjectId)) { + try (InputStream current = sharedStorage.download(sqlObjectId)) { // remove UTF-8 BOM if exists current.mark(3); byte[] byteSql = new byte[3]; @@ -427,14 +427,17 @@ private void upload() { } private String uploadToOSS(String filePath) { + if (null == context) { + throw new IllegalStateException("task not init, context is null"); + } // Public cloud scenario, need to upload files to OSS - CloudObjectStorageService cloudObjectStorageService = getCloudObjectStorageService(); - if (Objects.nonNull(cloudObjectStorageService) && cloudObjectStorageService.supported()) { + SharedStorage sharedStorage = context.getSharedStorage(); + if (sharedStorage.available()) { File file = new File(filePath); String ossAddress; try { - String objectName = cloudObjectStorageService.upload(file.getName(), file); - ossAddress = String.valueOf(cloudObjectStorageService.generateDownloadUrl(objectName)); + String objectName = sharedStorage.upload(file.getName(), file); + ossAddress = String.valueOf(sharedStorage.getDownloadURL(objectName)); log.info("Upload sql plan task result to cloud object storage successfully, objectName={}", objectName); } catch (Exception exception) { throw new RuntimeException(String.format( diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/dummy/DummyTask.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/dummy/DummyTask.java index d3f45f677e..c2be67a9da 100644 --- a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/dummy/DummyTask.java +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/dummy/DummyTask.java @@ -18,9 +18,9 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +import com.oceanbase.odc.service.task.TaskContext; import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.JobContext; -import com.oceanbase.odc.service.task.executor.task.TaskContext; import lombok.extern.slf4j.Slf4j; diff --git a/server/odc-service/src/test/java/com/oceanbase/odc/service/task/executor/task/BaseTaskTest.java b/server/odc-service/src/test/java/com/oceanbase/odc/service/task/executor/task/BaseTaskTest.java index d0f2e36690..857c920a40 100644 --- a/server/odc-service/src/test/java/com/oceanbase/odc/service/task/executor/task/BaseTaskTest.java +++ b/server/odc-service/src/test/java/com/oceanbase/odc/service/task/executor/task/BaseTaskTest.java @@ -24,10 +24,15 @@ import org.mockito.Mockito; import com.oceanbase.odc.common.util.SystemUtils; +import com.oceanbase.odc.service.task.ExceptionListener; +import com.oceanbase.odc.service.task.TaskContext; +import com.oceanbase.odc.service.task.TaskEventListener; +import com.oceanbase.odc.service.task.base.BaseTask; import com.oceanbase.odc.service.task.caller.DefaultJobContext; import com.oceanbase.odc.service.task.caller.JobContext; import com.oceanbase.odc.service.task.constants.JobEnvKeyConstants; -import com.oceanbase.odc.service.task.executor.server.TaskMonitor; +import com.oceanbase.odc.service.task.executor.DefaultTaskResult; +import com.oceanbase.odc.service.task.executor.DefaultTaskResultBuilder; import com.oceanbase.odc.service.task.schedule.JobIdentity; /** @@ -66,6 +71,11 @@ public ExceptionListener getExceptionListener() { public JobContext getJobContext() { return jobContext; } + + @Override + public TaskEventListener getTaskEventListener() { + return Mockito.mock(TaskEventListener.class); + } }); DefaultTaskResult taskResult = DefaultTaskResultBuilder.build(dummyBaseTask); DefaultTaskResultBuilder.assignErrorMessage(taskResult, dummyBaseTask); @@ -90,6 +100,11 @@ public ExceptionListener getExceptionListener() { public JobContext getJobContext() { return jobContext; } + + @Override + public TaskEventListener getTaskEventListener() { + return Mockito.mock(TaskEventListener.class); + } }); DefaultTaskResult taskResult = DefaultTaskResultBuilder.build(dummyBaseTask); DefaultTaskResultBuilder.assignErrorMessage(taskResult, dummyBaseTask); @@ -115,9 +130,6 @@ protected boolean doStart(JobContext context, TaskContext taskContext) throws Ex return true; } - protected TaskMonitor createTaskMonitor() { - return Mockito.mock(TaskMonitor.class); - } @Override protected void doStop() throws Exception {} From 1f8c0193680b3762ee7240a235cc278686a5accf Mon Sep 17 00:00:00 2001 From: "longpeng.zlp" Date: Tue, 29 Oct 2024 18:16:21 +0800 Subject: [PATCH 2/2] introduce task supervisor component --- .../agent/runtime/ExecutorRequestHandler.java | 3 +- .../service/task/supervisor/PortDetector.java | 111 ++++++++++++++++++ .../task/supervisor/TaskSupervisor.java | 74 ++++++++++++ .../task/supervisor/TaskSupervisorProxy.java | 70 +++++++++++ .../supervisor/endpoint/ExecutorEndpoint.java | 31 +++++ .../endpoint/SupervisorEndpoint.java | 30 +++++ .../service/task/supervisor/package-info.java | 25 ++++ .../task/supervisor/protocol/CommandType.java | 29 +++++ .../task/supervisor/protocol/TaskCommand.java | 61 ++++++++++ .../protocol/TaskCommandDeserializer.java | 50 ++++++++ .../protocol/TaskCommandExecutor.java | 66 +++++++++++ .../protocol/TaskCommandSender.java | 35 ++++++ .../protocol/TaskSupervisorServer.java | 30 +++++ .../protocol/command/GeneralTaskCommand.java | 68 +++++++++++ .../protocol/command/StartTaskCommand.java | 62 ++++++++++ .../proxy/LocalTaskSupervisorProxy.java | 105 +++++++++++++++++ .../proxy/RemoteTaskSupervisorProxy.java | 67 +++++++++++ 17 files changed, 916 insertions(+), 1 deletion(-) create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/PortDetector.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/TaskSupervisor.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/TaskSupervisorProxy.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/endpoint/ExecutorEndpoint.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/endpoint/SupervisorEndpoint.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/package-info.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/CommandType.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommand.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandDeserializer.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandExecutor.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandSender.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskSupervisorServer.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/command/GeneralTaskCommand.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/command/StartTaskCommand.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/proxy/LocalTaskSupervisorProxy.java create mode 100644 server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/proxy/RemoteTaskSupervisorProxy.java diff --git a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ExecutorRequestHandler.java b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ExecutorRequestHandler.java index 7dcc8cca35..49c29ef74b 100644 --- a/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ExecutorRequestHandler.java +++ b/server/odc-server/src/main/java/com/oceanbase/odc/agent/runtime/ExecutorRequestHandler.java @@ -57,7 +57,8 @@ public ExecutorRequestHandler() { this.executorBiz = new LogBizImpl(); } - public SuccessResponse process(HttpMethod httpMethod, String uri, String requestData) { + public SuccessResponse + process(HttpMethod httpMethod, String uri, String requestData) { if (uri == null || uri.trim().isEmpty()) { return Responses.single("request error: uri is empty."); } diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/PortDetector.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/PortDetector.java new file mode 100644 index 0000000000..d5e564385f --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/PortDetector.java @@ -0,0 +1,111 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor; + +import java.net.ServerSocket; +import java.util.HashSet; +import java.util.PriorityQueue; +import java.util.Set; + +import lombok.Setter; + +/** + * detect unused port for process in range + * @author longpeng.zlp + * @date 2024/10/25 16:25 + */ +public class PortDetector { + private static final PortDetector PORT_DETECTOR = new PortDetector(); + @Setter + private int maxPort; + @Setter + private int minPort; + @Setter + private int expiredMs; + private PriorityQueue allocatedPortInfos = new PriorityQueue<>((ap1, ap2) -> { + return Long.compare(ap1.getAllocatedTime(), ap2.getAllocatedTime()); + }); + private Set allocatedSets = new HashSet<>(); + + private PortDetector() { + maxPort = 65534; + minPort = 1000; + // 20s + expiredMs = 10000; + } + + public synchronized int getPort() { + long currentTimeMS = System.currentTimeMillis(); + // expire allocated tasks port + while (!allocatedPortInfos.isEmpty() && (currentTimeMS - allocatedPortInfos.peek().getAllocatedTime()) > expiredMs) { + AllocatedPortInfo allocatedPortInfo = allocatedPortInfos.poll(); + allocatedSets.remove(allocatedPortInfo.getPort()); + } + // go through and find available port + for (int i = minPort; i <= maxPort; ++i) { + if (allocatedSets.contains(i)) { + continue; + } + if (portInUse(i)) { + continue; + } + allocatedSets.add(i); + allocatedPortInfos.add(new AllocatedPortInfo(i, System.currentTimeMillis())); + return i; + } + throw new RuntimeException("port allocate failed"); + } + + private boolean portInUse(int port) { + ServerSocket socketServer = null; + try { + socketServer = new ServerSocket(port); + } catch (Throwable e) { + return true; + } finally { + if (null != socketServer) { + try { + socketServer.close(); + } catch (Throwable e) { + } + } + } + return false; + } + + public static PortDetector getInstance() { + return PORT_DETECTOR; + } + + private static final class AllocatedPortInfo { + private final int port; + private final long allocatedTime; + + private AllocatedPortInfo(int port, long allocatedTime) { + this.port = port; + this.allocatedTime = allocatedTime; + } + + public int getPort() { + return port; + } + + public long getAllocatedTime() { + return allocatedTime; + } + } +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/TaskSupervisor.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/TaskSupervisor.java new file mode 100644 index 0000000000..5872bc60e2 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/TaskSupervisor.java @@ -0,0 +1,74 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor; + +import com.oceanbase.odc.service.task.caller.JobContext; +import com.oceanbase.odc.service.task.caller.ProcessConfig; +import com.oceanbase.odc.service.task.supervisor.endpoint.ExecutorEndpoint; + +/** + * submit a task context and return a executor endpoint in local mode + * @author longpeng.zlp + * @date 2024/10/28 10:55 + */ +public class TaskSupervisor { + /** + * start task with given parameters + * @param jobContext + * @param processConfig + * @return + */ + public ExecutorEndpoint startTask(JobContext jobContext, ProcessConfig processConfig) { + return null; + } + + /** + * stop task + * @param jobContext + */ + public boolean stopTask(ExecutorEndpoint executorEndpoint, JobContext jobContext) { + return true; + } + + /** + * modify a task without stop it + * @param executorEndpoint + * @param jobContext + */ + public boolean modifyTask(ExecutorEndpoint executorEndpoint, JobContext jobContext) { + return true; + } + + /** + * finish a task, it will not be started again + * @param executorEndpoint + * @param jobContext + */ + public boolean finishTask(ExecutorEndpoint executorEndpoint, JobContext jobContext) { + return true; + } + + /** + * compatible with old job caller + * @param executorEndpoint + * @param jobContext + * @return + */ + public boolean canBeFinished(ExecutorEndpoint executorEndpoint, JobContext jobContext) { + return true; + } +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/TaskSupervisorProxy.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/TaskSupervisorProxy.java new file mode 100644 index 0000000000..b62254b6da --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/TaskSupervisorProxy.java @@ -0,0 +1,70 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor; + +import com.oceanbase.odc.service.task.caller.JobContext; +import com.oceanbase.odc.service.task.caller.ProcessConfig; +import com.oceanbase.odc.service.task.supervisor.endpoint.ExecutorEndpoint; +import com.oceanbase.odc.service.task.supervisor.endpoint.SupervisorEndpoint; + +/** + * execute remote/local call for given supervisor endpoint + * @author longpeng.zlp + * @date 2024/10/29 11:48 + */ +public interface TaskSupervisorProxy { + /** + * execute start task command to supervisorEndpoint + * @param supervisorEndpoint + * @param jobContext + * @param processConfig + * @return + */ + ExecutorEndpoint startTask(SupervisorEndpoint supervisorEndpoint, JobContext jobContext, ProcessConfig processConfig); + + /** + * execute stop task command to supervisorEndpoint + * @param supervisorEndpoint + * @param jobContext + * @return + */ + boolean stopTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, JobContext jobContext); + + /** + * execute modify task command to supervisorEndpoint + * @param supervisorEndpoint + * @param jobContext + * @return + */ + boolean modifyTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, JobContext jobContext); + + /** + * execute finish task command to supervisorEndpoint + * @param supervisorEndpoint + * @param jobContext + * @return + */ + boolean finishTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, JobContext jobContext); + + /** + * detect can be finish command to supervisorEndpoint + * @param supervisorEndpoint + * @param jobContext + * @return + */ + boolean canBeFinished(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, JobContext jobContext); +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/endpoint/ExecutorEndpoint.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/endpoint/ExecutorEndpoint.java new file mode 100644 index 0000000000..4fb5a92c7c --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/endpoint/ExecutorEndpoint.java @@ -0,0 +1,31 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.endpoint; + +import lombok.Data; + +/** + * @author longpeng.zlp + * @date 2024/10/28 16:41 + */ +@Data +public class ExecutorEndpoint { + private final String protocol; + private final String host; + private final String port; + private final String identifier; +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/endpoint/SupervisorEndpoint.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/endpoint/SupervisorEndpoint.java new file mode 100644 index 0000000000..ac2382677b --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/endpoint/SupervisorEndpoint.java @@ -0,0 +1,30 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.endpoint; + +import lombok.Data; + +/** + * @author longpeng.zlp + * @date 2024/10/29 14:24 + */ +@Data +public class SupervisorEndpoint { + public static final SupervisorEndpoint SELF_ENDPOINT = new SupervisorEndpoint("localhost", "-1"); + private final String host; + private final String port; +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/package-info.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/package-info.java new file mode 100644 index 0000000000..b63b5a4751 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/package-info.java @@ -0,0 +1,25 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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. + */ + +/** + * @author longpeng.zlp + * @date 2024/10/28 16:37 + */ +@Evolving +package com.oceanbase.odc.service.task.supervisor; + +import org.apache.hadoop.classification.InterfaceStability.Evolving; + diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/CommandType.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/CommandType.java new file mode 100644 index 0000000000..511c9f0fc5 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/CommandType.java @@ -0,0 +1,29 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.protocol; + +/** + * @author longpeng.zlp + * @date 2024/10/29 15:21 + */ +public enum CommandType { + START, + STOP, + MODIFY, + FINISH, + CAN_BE_FINISHED +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommand.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommand.java new file mode 100644 index 0000000000..bd19ef2164 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommand.java @@ -0,0 +1,61 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.protocol; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.oceanbase.odc.common.json.JsonUtils; +import com.oceanbase.odc.service.task.caller.JobContext; + +import lombok.Getter; + +/** + * @author longpeng.zlp + * @date 2024/10/29 15:13 + */ +public abstract class TaskCommand { + protected static final ObjectMapper objectMapper = new ObjectMapper(); + protected static final int COMMAND_VERSION = 0; + protected static final String VERSION_NAME = "version"; + protected static final String JOB_CONTEXT_NAME = "jobContext"; + protected static final String COMMAND_TYPE_NAME = "command"; + @Getter + protected JobContext jobContext; + @Getter + protected int version; + + public abstract CommandType commandType(); + + public String serialize() { + ObjectNode jsonNode = objectMapper.createObjectNode(); + jsonNode.put(VERSION_NAME, version); + jsonNode.put(JOB_CONTEXT_NAME, JsonUtils.toJson(jobContext)); + jsonNode.put(COMMAND_TYPE_NAME, commandType().toString()); + append(jsonNode); + return jsonNode.toPrettyString(); + } + + protected void fillCommonFields(JsonNode jsonNode) { + JsonNode versionNode = jsonNode.get(VERSION_NAME); + version = null == versionNode ? 0 : versionNode.asInt(); + JsonNode jobContextNode = jsonNode.get(JOB_CONTEXT_NAME); + jobContext = jobContextNode == null ? null : JsonUtils.fromJson(jobContextNode.asText(), JobContext.class); + } + + public abstract void append(ObjectNode objectNode); +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandDeserializer.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandDeserializer.java new file mode 100644 index 0000000000..a5876ab7e0 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandDeserializer.java @@ -0,0 +1,50 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.protocol; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.oceanbase.odc.service.task.supervisor.protocol.command.GeneralTaskCommand; +import com.oceanbase.odc.service.task.supervisor.protocol.command.StartTaskCommand; + +/** + * @author longpeng.zlp + * @date 2024/10/29 15:59 + */ +public class TaskCommandDeserializer { + private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); + public TaskCommand deserializeTaskCommand(String commandStr) throws JsonProcessingException { + JsonNode objectNode = OBJECT_MAPPER.readTree(commandStr); + JsonNode commandTypeNode = objectNode.get(TaskCommand.COMMAND_TYPE_NAME); + if (null == commandTypeNode) { + throw new IllegalStateException("invalid command, str=" + commandStr); + } + CommandType commandType = CommandType.valueOf(commandTypeNode.asText()); + switch (commandType) { + case START: + return StartTaskCommand.fromJsonNode(objectNode); + case FINISH: + case MODIFY: + case STOP: + case CAN_BE_FINISHED: + return GeneralTaskCommand.fromJsonNode(GeneralTaskCommand::new, objectNode, commandType); + default: + throw new IllegalStateException("not support command type, str=" + commandType); + } + } +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandExecutor.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandExecutor.java new file mode 100644 index 0000000000..f692f5f973 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandExecutor.java @@ -0,0 +1,66 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.protocol; + +import java.util.function.BiFunction; +import java.util.function.Consumer; + +import com.oceanbase.odc.common.json.JsonUtils; +import com.oceanbase.odc.service.task.caller.JobContext; +import com.oceanbase.odc.service.task.supervisor.TaskSupervisor; +import com.oceanbase.odc.service.task.supervisor.endpoint.ExecutorEndpoint; +import com.oceanbase.odc.service.task.supervisor.protocol.command.GeneralTaskCommand; +import com.oceanbase.odc.service.task.supervisor.protocol.command.StartTaskCommand; + +/** + * @author longpeng.zlp + * @date 2024/10/29 17:45 + */ +public class TaskCommandExecutor { + public TaskSupervisor taskSupervisor; + + public void onCommand(TaskCommand taskCommand, Consumer responseConsumer) { + String ret = null; + switch (taskCommand.commandType()) { + case START: + StartTaskCommand startTaskCommand = (StartTaskCommand) taskCommand; + ExecutorEndpoint endpoint = taskSupervisor.startTask(startTaskCommand.getJobContext(), startTaskCommand.getProcessConfig()); + ret = JsonUtils.toJson(endpoint); + break; + default: + boolean succeed = callTaskSupervisorFunc((GeneralTaskCommand) taskCommand); + ret = String.valueOf(succeed); + break; + } + responseConsumer.accept(ret); + } + + protected boolean callTaskSupervisorFunc(GeneralTaskCommand generalTaskCommand) { + switch (generalTaskCommand.commandType()) { + case STOP: + return taskSupervisor.stopTask(generalTaskCommand.getExecutorEndpoint(), generalTaskCommand.getJobContext()); + case MODIFY: + return taskSupervisor.modifyTask(generalTaskCommand.getExecutorEndpoint(), generalTaskCommand.getJobContext()); + case FINISH: + return taskSupervisor.finishTask(generalTaskCommand.getExecutorEndpoint(), generalTaskCommand.getJobContext()); + case CAN_BE_FINISHED: + return taskSupervisor.canBeFinished(generalTaskCommand.getExecutorEndpoint(), generalTaskCommand.getJobContext()); + default: + throw new IllegalStateException("not recognized command " + generalTaskCommand); + } + } +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandSender.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandSender.java new file mode 100644 index 0000000000..767018d6a3 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskCommandSender.java @@ -0,0 +1,35 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.protocol; + +import com.oceanbase.odc.service.task.supervisor.endpoint.SupervisorEndpoint; + +/** + * @author longpeng.zlp + * @date 2024/10/29 15:01 + */ +public class TaskCommandSender { + /** + * send command to supervisor end point and return reponse + * @param supervisorEndpoint + * @param taskCommand + * @return + */ + public String sendCommand(SupervisorEndpoint supervisorEndpoint, TaskCommand taskCommand) { + return ""; + } +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskSupervisorServer.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskSupervisorServer.java new file mode 100644 index 0000000000..53222172a5 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/TaskSupervisorServer.java @@ -0,0 +1,30 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.protocol; + +/** + * + * @author longpeng.zlp + * @date 2024/10/29 18:05 + */ +public class TaskSupervisorServer { + private TaskCommandExecutor taskCommandExecutor; + + /** + * impl http protocol here + */ +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/command/GeneralTaskCommand.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/command/GeneralTaskCommand.java new file mode 100644 index 0000000000..956753632f --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/command/GeneralTaskCommand.java @@ -0,0 +1,68 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.protocol.command; + +import java.util.function.Supplier; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.oceanbase.odc.common.json.JsonUtils; +import com.oceanbase.odc.service.task.caller.JobContext; +import com.oceanbase.odc.service.task.supervisor.endpoint.ExecutorEndpoint; +import com.oceanbase.odc.service.task.supervisor.protocol.CommandType; +import com.oceanbase.odc.service.task.supervisor.protocol.TaskCommand; + +import lombok.Getter; + +/** + * @author longpeng.zlp + * @date 2024/10/29 15:47 + */ +public class GeneralTaskCommand extends TaskCommand { + protected static final String EXECUTOR_ENC_POINT_STR = "executorEndpoint"; + @Getter + protected ExecutorEndpoint executorEndpoint; + protected CommandType commandType; + + public void append(ObjectNode objectNode) { + objectNode.put(EXECUTOR_ENC_POINT_STR, JsonUtils.toJson(executorEndpoint)); + } + + public CommandType commandType() { + return commandType; + } + + public static T fromJsonNode(Supplier commandSupplier, JsonNode jsonNode, CommandType commandType) { + T command = commandSupplier.get(); + JsonNode executorEndpointNode = jsonNode.get(EXECUTOR_ENC_POINT_STR); + ExecutorEndpoint endpoint = JsonUtils.fromJson(null == executorEndpointNode ? null : executorEndpointNode.asText(), ExecutorEndpoint.class); + command.fillCommonFields(jsonNode); + command.executorEndpoint = endpoint; + command.commandType = commandType; + return command; + } + + public static GeneralTaskCommand create(JobContext jobContext, ExecutorEndpoint endpoint, CommandType commandType) { + GeneralTaskCommand ret = new GeneralTaskCommand(); + ret.commandType = commandType; + ret.version = COMMAND_VERSION; + ret.jobContext = jobContext; + ret.executorEndpoint = endpoint; + return ret; + } + +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/command/StartTaskCommand.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/command/StartTaskCommand.java new file mode 100644 index 0000000000..f58573b16d --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/protocol/command/StartTaskCommand.java @@ -0,0 +1,62 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.protocol.command; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; +import com.oceanbase.odc.common.json.JsonUtils; +import com.oceanbase.odc.service.task.caller.JobContext; +import com.oceanbase.odc.service.task.caller.ProcessConfig; +import com.oceanbase.odc.service.task.supervisor.protocol.CommandType; +import com.oceanbase.odc.service.task.supervisor.protocol.TaskCommand; + +import lombok.Getter; + +/** + * @author longpeng.zlp + * @date 2024/10/29 15:25 + */ +public class StartTaskCommand extends TaskCommand { + protected static final String PROCESS_CONFIG_STR = "processConfig"; + @Getter + private ProcessConfig processConfig; + @Override + public CommandType commandType() { + return CommandType.START; + } + + @Override + public void append(ObjectNode objectNode) { + objectNode.put(PROCESS_CONFIG_STR, JsonUtils.toJson(processConfig)); + } + + public static StartTaskCommand fromJsonNode(JsonNode jsonNode) { + StartTaskCommand startTaskCommand = new StartTaskCommand(); + JsonNode processConfigNode = jsonNode.get(PROCESS_CONFIG_STR); + startTaskCommand.processConfig = JsonUtils.fromJson(null == processConfigNode ? null : processConfigNode.asText(), ProcessConfig.class); + startTaskCommand.fillCommonFields(jsonNode); + return startTaskCommand; + } + + public static StartTaskCommand create(JobContext jobContext, ProcessConfig processConfig) { + StartTaskCommand ret = new StartTaskCommand(); + ret.processConfig = processConfig; + ret.jobContext = jobContext; + ret.version = COMMAND_VERSION; + return ret; + } +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/proxy/LocalTaskSupervisorProxy.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/proxy/LocalTaskSupervisorProxy.java new file mode 100644 index 0000000000..6e6f73b679 --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/proxy/LocalTaskSupervisorProxy.java @@ -0,0 +1,105 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.proxy; + +import com.oceanbase.odc.service.task.caller.JobContext; +import com.oceanbase.odc.service.task.caller.ProcessConfig; +import com.oceanbase.odc.service.task.supervisor.TaskSupervisor; +import com.oceanbase.odc.service.task.supervisor.TaskSupervisorProxy; +import com.oceanbase.odc.service.task.supervisor.endpoint.ExecutorEndpoint; +import com.oceanbase.odc.service.task.supervisor.endpoint.SupervisorEndpoint; + +import lombok.extern.slf4j.Slf4j; + +/** + * proxy to route task command local impl + * @author longpeng.zlp + * @date 2024/10/29 11:43 + */ +@Slf4j +public class LocalTaskSupervisorProxy implements TaskSupervisorProxy{ + private TaskSupervisor localTaskSupervisor; + private RemoteTaskSupervisorProxy remoteTaskSupervisorProxy; + private SupervisorEndpoint localEndPoint; + + @Override + public ExecutorEndpoint startTask(SupervisorEndpoint supervisorEndpoint, JobContext jobContext, + ProcessConfig processConfig) { + if (isLocalCommandCall(supervisorEndpoint)) { + log.info("local call start task, supervisorEndpoint={}, jobContext={}, processConfig={}", supervisorEndpoint, jobContext, processConfig); + return localTaskSupervisor.startTask(jobContext, processConfig); + } else { + log.info("remote call start task, supervisorEndpoint={}, jobContext={}, processConfig={}", supervisorEndpoint, jobContext, processConfig); + return remoteTaskSupervisorProxy.startTask(supervisorEndpoint, jobContext, processConfig); + } + } + + @Override + public boolean stopTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, + JobContext jobContext) { + if (isLocalCommandCall(supervisorEndpoint)) { + log.info("local call stop task, supervisorEndpoint={}, executorEndpoint={}, jobContext={}", supervisorEndpoint, executorEndpoint, jobContext); + return localTaskSupervisor.stopTask(executorEndpoint, jobContext); + } else { + log.info("remote call stop task, supervisorEndpoint={}, executorEndpoint={}, jobContext={}", supervisorEndpoint, executorEndpoint, jobContext); + return remoteTaskSupervisorProxy.stopTask(supervisorEndpoint, executorEndpoint, jobContext); + } + } + + @Override + public boolean modifyTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, + JobContext jobContext) { + if (isLocalCommandCall(supervisorEndpoint)) { + log.info("local call modify task, supervisorEndpoint={}, executorEndpoint={}, jobContext={}", supervisorEndpoint, executorEndpoint, jobContext); + return localTaskSupervisor.modifyTask(executorEndpoint, jobContext); + } else { + log.info("remote call modify task, supervisorEndpoint={}, executorEndpoint={}, jobContext={}", supervisorEndpoint, executorEndpoint, jobContext); + return remoteTaskSupervisorProxy.modifyTask(supervisorEndpoint, executorEndpoint, jobContext); + } + } + + @Override + public boolean finishTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, + JobContext jobContext) { + if (isLocalCommandCall(supervisorEndpoint)) { + log.info("local call finish task, supervisorEndpoint={}, executorEndpoint={}, jobContext={}", supervisorEndpoint, executorEndpoint, jobContext); + return localTaskSupervisor.finishTask(executorEndpoint, jobContext); + } else { + log.info("remote call finish task, supervisorEndpoint={}, executorEndpoint={}, jobContext={}", supervisorEndpoint, executorEndpoint, jobContext); + return remoteTaskSupervisorProxy.finishTask(supervisorEndpoint, executorEndpoint, jobContext); + } + } + + @Override + public boolean canBeFinished(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, + JobContext jobContext) { + if (isLocalCommandCall(supervisorEndpoint)) { + log.info("local call canBeFinished task, supervisorEndpoint={}, executorEndpoint={}, jobContext={}", supervisorEndpoint, executorEndpoint, jobContext); + return localTaskSupervisor.canBeFinished(executorEndpoint, jobContext); + } else { + log.info("remote call canBeFinished task, supervisorEndpoint={}, executorEndpoint={}, jobContext={}", supervisorEndpoint, executorEndpoint, jobContext); + return remoteTaskSupervisorProxy.canBeFinished(supervisorEndpoint, executorEndpoint, jobContext); + } + } + + protected boolean isLocalCommandCall(SupervisorEndpoint supervisorEndpoint) { + if (null == supervisorEndpoint) { + throw new IllegalStateException("end point must be given for task supervisor proxy"); + } + return supervisorEndpoint.equals(SupervisorEndpoint.SELF_ENDPOINT) || supervisorEndpoint.equals(localEndPoint); + } +} diff --git a/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/proxy/RemoteTaskSupervisorProxy.java b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/proxy/RemoteTaskSupervisorProxy.java new file mode 100644 index 0000000000..993b38728c --- /dev/null +++ b/server/odc-service/src/main/java/com/oceanbase/odc/service/task/supervisor/proxy/RemoteTaskSupervisorProxy.java @@ -0,0 +1,67 @@ +/* + * Copyright (c) 2024 OceanBase. + * + * Licensed 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 com.oceanbase.odc.service.task.supervisor.proxy; + +import com.oceanbase.odc.common.json.JsonUtils; +import com.oceanbase.odc.service.task.caller.JobContext; +import com.oceanbase.odc.service.task.caller.ProcessConfig; +import com.oceanbase.odc.service.task.supervisor.TaskSupervisorProxy; +import com.oceanbase.odc.service.task.supervisor.endpoint.ExecutorEndpoint; +import com.oceanbase.odc.service.task.supervisor.endpoint.SupervisorEndpoint; +import com.oceanbase.odc.service.task.supervisor.protocol.CommandType; +import com.oceanbase.odc.service.task.supervisor.protocol.TaskCommandSender; +import com.oceanbase.odc.service.task.supervisor.protocol.command.GeneralTaskCommand; +import com.oceanbase.odc.service.task.supervisor.protocol.command.StartTaskCommand; + +/** + * @author longpeng.zlp + * @date 2024/10/29 14:42 + */ +public class RemoteTaskSupervisorProxy implements TaskSupervisorProxy { + private TaskCommandSender taskCommandSender; + @Override + public ExecutorEndpoint startTask(SupervisorEndpoint supervisorEndpoint, JobContext jobContext, + ProcessConfig processConfig) { + return JsonUtils.fromJson(taskCommandSender.sendCommand(supervisorEndpoint, StartTaskCommand.create(jobContext, processConfig)), ExecutorEndpoint.class); + } + + @Override + public boolean stopTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, + JobContext jobContext) { + return Boolean.valueOf(taskCommandSender.sendCommand(supervisorEndpoint, GeneralTaskCommand.create(jobContext, executorEndpoint, CommandType.STOP))); + } + + @Override + public boolean modifyTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, + JobContext jobContext) { + return Boolean.valueOf(taskCommandSender.sendCommand(supervisorEndpoint, GeneralTaskCommand.create(jobContext, executorEndpoint, CommandType.MODIFY))); + } + + @Override + public boolean finishTask(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, + JobContext jobContext) { + return Boolean.valueOf(taskCommandSender.sendCommand(supervisorEndpoint, GeneralTaskCommand.create(jobContext, executorEndpoint, CommandType.FINISH))); + } + + @Override + public boolean canBeFinished(SupervisorEndpoint supervisorEndpoint, ExecutorEndpoint executorEndpoint, + JobContext jobContext) { + return Boolean.valueOf(taskCommandSender.sendCommand(supervisorEndpoint, GeneralTaskCommand.create(jobContext, executorEndpoint, CommandType.CAN_BE_FINISHED))); + } + + +}