This is an automated email from the ASF dual-hosted git repository. technoboy pushed a commit to branch refactor-worker in repository https://gitbox.apache.org/repos/asf/incubator-dolphinscheduler.git
commit 0e0c47925bef68dafbed5df239a8126932741505 Merge: ff86dc7 63b76d7 Author: Tboy <[email protected]> AuthorDate: Fri Feb 21 20:32:58 2020 +0800 Merge branch 'refactor-worker' into refactor-worker .../dolphinscheduler/server/master/runner/MasterBaseTaskExecThread.java | 2 -- .../dolphinscheduler/server/worker/runner/TaskScheduleThread.java | 2 +- 2 files changed, 1 insertion(+), 3 deletions(-) diff --cc dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/MasterBaseTaskExecThread.java index 4ae057a,a261b34..09005a1 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/MasterBaseTaskExecThread.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/MasterBaseTaskExecThread.java @@@ -128,23 -124,11 +128,21 @@@ public class MasterBaseTaskExecThread i // TODO send task to worker public void sendToWorker(TaskInstance taskInstance){ final Address address = new Address("127.0.0.1", 12346); - - ExecuteTaskRequestCommand taskRequestCommand = new ExecuteTaskRequestCommand(FastJsonSerializer.serializeToString(taskInstance)); + /** + * set taskInstance relation + */ + TaskInstance destTaskInstance = setTaskInstanceRelation(taskInstance); + + ExecuteTaskRequestCommand taskRequestCommand = new ExecuteTaskRequestCommand( + FastJsonSerializer.serializeToString(convertToTaskInfo(destTaskInstance))); try { - Command responseCommand = nettyRemotingClient.sendSync(address, taskRequestCommand.convert2Command(), Integer.MAX_VALUE); - ExecuteTaskAckCommand taskAckCommand = FastJsonSerializer.deserialize(responseCommand.getBody(), ExecuteTaskAckCommand.class); + Command responseCommand = nettyRemotingClient.sendSync(address, + taskRequestCommand.convert2Command(), Integer.MAX_VALUE); + + ExecuteTaskAckCommand taskAckCommand = FastJsonSerializer.deserialize( + responseCommand.getBody(), ExecuteTaskAckCommand.class); + logger.info("taskAckCommand : {}",taskAckCommand); - processService.changeTaskState(ExecutionStatus.of(taskAckCommand.getStatus()), taskAckCommand.getStartTime(), taskAckCommand.getHost(),
