Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -968,7 +968,7 @@ public void childEvent(CuratorFramework client, PathChildrenCacheEvent event)
announcement.getTaskType(),
zkWorker.getWorker(),
TaskLocation.unknown(),
runningTasks.get(taskId).getDataSource()
announcement.getTaskDataSource()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If the task has already been shutdown, will this make it start back up again?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I guess there's already a "we'll add it if we don't know about it" thing immediately below, so the check must be elsewhere.

);
final RemoteTaskRunnerWorkItem existingItem = runningTasks.putIfAbsent(
taskId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -387,7 +387,8 @@ public void fullSync(List<WorkerHistoryItem> changes)
announcement.getTaskType(),
announcement.getTaskResource(),
TaskStatus.failure(announcement.getTaskId()),
announcement.getTaskLocation()
announcement.getTaskLocation(),
announcement.getTaskDataSource()
));
}
}
Expand Down Expand Up @@ -423,7 +424,8 @@ public void deltaSync(List<WorkerHistoryItem> changes)
announcement.getTaskType(),
announcement.getTaskResource(),
TaskStatus.failure(announcement.getTaskId()),
announcement.getTaskLocation()
announcement.getTaskLocation(),
announcement.getTaskDataSource()
));
}
} else if (change instanceof WorkerHistoryItem.Metadata) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@
import io.druid.indexing.common.task.Task;
import io.druid.indexing.common.task.TaskResource;

import javax.annotation.Nullable;

/**
* Used by workers to announce the status of tasks they are currently running. This class is immutable.
*/
Expand All @@ -38,21 +40,25 @@ public class TaskAnnouncement
private final TaskResource taskResource;
private final TaskLocation taskLocation;

@Nullable
private final String taskDataSource; // nullable for backward compatibility

public static TaskAnnouncement create(Task task, TaskStatus status, TaskLocation location)
{
return create(task.getId(), task.getType(), task.getTaskResource(), status, location);
return create(task.getId(), task.getType(), task.getTaskResource(), status, location, task.getDataSource());
}

public static TaskAnnouncement create(
String taskId,
String taskType,
TaskResource resource,
TaskStatus status,
TaskLocation location
TaskLocation location,
String taskDataSource
)
{
Preconditions.checkArgument(status.getId().equals(taskId), "task id == status id");
return new TaskAnnouncement(null, taskType, null, status, resource, location);
return new TaskAnnouncement(null, taskType, null, status, resource, location, taskDataSource);
}

@JsonCreator
Expand All @@ -62,7 +68,8 @@ private TaskAnnouncement(
@JsonProperty("status") TaskState status,
@JsonProperty("taskStatus") TaskStatus taskStatus,
@JsonProperty("taskResource") TaskResource taskResource,
@JsonProperty("taskLocation") TaskLocation taskLocation
@JsonProperty("taskLocation") TaskLocation taskLocation,
@JsonProperty("taskDataSource") String taskDataSource
)
{
this.taskType = taskType;
Expand All @@ -74,6 +81,7 @@ private TaskAnnouncement(
}
this.taskResource = taskResource == null ? new TaskResource(this.taskStatus.getId(), 1) : taskResource;
this.taskLocation = taskLocation == null ? TaskLocation.unknown() : taskLocation;
this.taskDataSource = taskDataSource;
}

@JsonProperty("id")
Expand Down Expand Up @@ -112,13 +120,21 @@ public TaskLocation getTaskLocation()
return taskLocation;
}

@JsonProperty("taskDataSource")
public String getTaskDataSource()
{
return taskDataSource;
}

@Override
public String toString()
{
return "TaskAnnouncement{" +
"taskStatus=" + taskStatus +
"taskType=" + taskType +
", taskStatus=" + taskStatus +
", taskResource=" + taskResource +
", taskLocation=" + taskLocation +
", taskDataSource=" + taskDataSource +
'}';
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,8 @@ private void cleanupStaleAnnouncements() throws Exception
announcement.getTaskType(),
announcement.getTaskResource(),
completionStatus,
TaskLocation.unknown()
TaskLocation.unknown(),
announcement.getTaskDataSource()
)
);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -447,6 +447,33 @@ public void testWorkerDisabled() throws Exception
Assert.assertEquals("", Iterables.getOnlyElement(remoteTaskRunner.getWorkers()).getWorker().getVersion());
}

@Test
public void testRestartRemoteTaskRunner() throws Exception
{
doSetup();
remoteTaskRunner.run(task);

Assert.assertTrue(taskAnnounced(task.getId()));
mockWorkerRunningTask(task);
Assert.assertTrue(workerRunningTask(task.getId()));

remoteTaskRunner.stop();
makeRemoteTaskRunner(new TestRemoteTaskRunnerConfig(new Period("PT5S")));
final RemoteTaskRunnerWorkItem newWorkItem = remoteTaskRunner
.getKnownTasks()
.stream()
.filter(workItem -> workItem.getTaskId().equals(task.getId()))
.findFirst()
.orElse(null);
final ListenableFuture<TaskStatus> result = newWorkItem.getResult();

mockWorkerCompleteSuccessfulTask(task);
Assert.assertTrue(workerCompletedTask(result));

Assert.assertEquals(task.getId(), result.get().getId());
Assert.assertEquals(TaskState.SUCCESS, result.get().getStatusCode());
}

private void doSetup() throws Exception
{
makeWorker();
Expand Down