mirror of https://github.com/apache/druid.git
address cr again
This commit is contained in:
parent
55db06ccb1
commit
4c23a5e9f6
|
@ -30,24 +30,15 @@ import java.util.Set;
|
|||
*/
|
||||
public class ImmutableZkWorker
|
||||
{
|
||||
private final ZkWorker mutableZkWorker;
|
||||
private final Worker worker;
|
||||
private final int currCapacityUsed;
|
||||
private final Set<String> availabilityGroups;
|
||||
|
||||
public ImmutableZkWorker(
|
||||
ZkWorker mutableZkWorker
|
||||
)
|
||||
public ImmutableZkWorker(Worker worker, int currCapacityUsed, Set<String> availabilityGroups)
|
||||
{
|
||||
this.mutableZkWorker = mutableZkWorker;
|
||||
this.worker = mutableZkWorker.getWorker();
|
||||
this.currCapacityUsed = mutableZkWorker.getCurrCapacityUsed();
|
||||
this.availabilityGroups = ImmutableSet.copyOf(mutableZkWorker.getAvailabilityGroups());
|
||||
}
|
||||
|
||||
public ZkWorker getMutableZkWorker()
|
||||
{
|
||||
return mutableZkWorker;
|
||||
this.worker = worker;
|
||||
this.currCapacityUsed = currCapacityUsed;
|
||||
this.availabilityGroups = ImmutableSet.copyOf(availabilityGroups);
|
||||
}
|
||||
|
||||
public Worker getWorker()
|
||||
|
|
|
@ -64,7 +64,6 @@ import org.apache.zookeeper.KeeperException;
|
|||
import org.jboss.netty.handler.codec.http.HttpResponseStatus;
|
||||
import org.joda.time.DateTime;
|
||||
|
||||
import javax.annotation.Nullable;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.net.MalformedURLException;
|
||||
|
@ -522,7 +521,7 @@ public class RemoteTaskRunner implements TaskRunner, TaskLogStreamer
|
|||
return true;
|
||||
} else {
|
||||
// Nothing running this task, announce it in ZK for a worker to run it
|
||||
Optional<ZkWorker> zkWorker = strategy.findWorkerForTask(
|
||||
final Optional<ImmutableZkWorker> immutableZkWorker = strategy.findWorkerForTask(
|
||||
ImmutableMap.copyOf(
|
||||
Maps.transformEntries(
|
||||
zkWorkers,
|
||||
|
@ -540,8 +539,9 @@ public class RemoteTaskRunner implements TaskRunner, TaskLogStreamer
|
|||
),
|
||||
task
|
||||
);
|
||||
if (zkWorker.isPresent()) {
|
||||
announceTask(task, zkWorker.get(), taskRunnerWorkItem);
|
||||
if (immutableZkWorker.isPresent()) {
|
||||
final ZkWorker zkWorker = zkWorkers.get(immutableZkWorker.get().getWorker().getHost());
|
||||
announceTask(task, zkWorker, taskRunnerWorkItem);
|
||||
return true;
|
||||
} else {
|
||||
return false;
|
||||
|
|
|
@ -44,9 +44,9 @@ public interface WorkerSelectStrategy
|
|||
* @param zkWorkers An immutable map of workers to choose from.
|
||||
* @param task The task to assign.
|
||||
*
|
||||
* @return A {@link io.druid.indexing.overlord.ZkWorker} to run the task if one is available.
|
||||
* @return A {@link io.druid.indexing.overlord.ImmutableZkWorker} to run the task if one is available.
|
||||
*/
|
||||
public Optional<ZkWorker> findWorkerForTask(
|
||||
public Optional<ImmutableZkWorker> findWorkerForTask(
|
||||
final ImmutableMap<String, ImmutableZkWorker> zkWorkers,
|
||||
final Task task
|
||||
);
|
||||
|
|
Loading…
Reference in New Issue