Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -20,8 +20,12 @@
package org.apache.druid.indexing.common.actions;

import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import org.apache.druid.indexing.common.task.IndexTaskUtils;
import org.apache.druid.indexing.common.task.Task;
import org.apache.druid.java.util.common.Stopwatch;
import org.apache.druid.java.util.common.jackson.JacksonUtils;
import org.apache.druid.java.util.emitter.EmittingLogger;
import org.apache.druid.java.util.emitter.service.ServiceMetricEvent;
Expand All @@ -34,6 +38,8 @@ public class LocalTaskActionClient implements TaskActionClient
{
private static final EmittingLogger log = new EmittingLogger(LocalTaskActionClient.class);

private static final String RUN_TIME_METRIC = "task/action/run/time";

private final Task task;
private final TaskActionToolbox toolbox;

Expand All @@ -52,10 +58,33 @@ public <RetType> RetType submit(TaskAction<RetType> taskAction)
log.debug("Performing action for task[%s]: %s", task.getId(), taskAction);
final long performStartTime = System.currentTimeMillis();
final RetType result = performAction(taskAction);
emitTimerMetric("task/action/run/time", taskAction, System.currentTimeMillis() - performStartTime);
emitTimerMetric(RUN_TIME_METRIC, taskAction, System.currentTimeMillis() - performStartTime);
return result;
}

@Override
public <RetType> ListenableFuture<RetType> submitAsync(TaskAction<RetType> taskAction)
{
try {
if (taskAction.canPerformAsync(task, toolbox)) {
final Stopwatch actionRunTime = Stopwatch.createStarted();
return Futures.transform(
taskAction.performAsync(task, toolbox),
v -> {
emitTimerMetric(RUN_TIME_METRIC, taskAction, actionRunTime.millisElapsed());
return v;
},
MoreExecutors.directExecutor()
);
} else {
return Futures.immediateFuture(submit(taskAction));
}
}
catch (Exception e) {
return Futures.immediateFailedFuture(e);
}
}

private <R> R performAction(TaskAction<R> taskAction)
{
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.core.type.TypeReference;
import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.ListenableFuture;
import org.apache.druid.error.DruidException;
import org.apache.druid.indexing.common.LockGranularity;
import org.apache.druid.indexing.common.TaskLockType;
Expand Down Expand Up @@ -50,7 +51,6 @@
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.concurrent.Future;
import java.util.concurrent.ThreadLocalRandom;

/**
Expand Down Expand Up @@ -194,7 +194,7 @@ public boolean canPerformAsync(Task task, TaskActionToolbox toolbox)
}

@Override
public Future<SegmentIdWithShardSpec> performAsync(Task task, TaskActionToolbox toolbox)
public ListenableFuture<SegmentIdWithShardSpec> performAsync(Task task, TaskActionToolbox toolbox)
{
if (!toolbox.canBatchSegmentAllocation()) {
throw new ISE("Batched segment allocation is disabled");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,8 @@

package org.apache.druid.indexing.common.actions;

import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture;
import com.google.inject.Inject;
import org.apache.druid.guice.ManageLifecycle;
import org.apache.druid.indexing.common.LockGranularity;
Expand Down Expand Up @@ -54,9 +56,7 @@
import java.util.Set;
import java.util.TreeSet;
import java.util.concurrent.BlockingDeque;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Future;
import java.util.concurrent.LinkedBlockingDeque;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
Expand Down Expand Up @@ -211,7 +211,7 @@ public int size()
* Queues a SegmentAllocateRequest. The returned future may complete successfully
* with a non-null value or with a non-null value.
*/
public Future<SegmentIdWithShardSpec> add(SegmentAllocateRequest request)
public ListenableFuture<SegmentIdWithShardSpec> add(SegmentAllocateRequest request)
{
if (!isLeader.get()) {
throw new ISE("Cannot allocate segment if not leader.");
Expand All @@ -220,7 +220,7 @@ public Future<SegmentIdWithShardSpec> add(SegmentAllocateRequest request)
}

final AllocateRequestKey requestKey = new AllocateRequestKey(request);
final AtomicReference<Future<SegmentIdWithShardSpec>> futureReference = new AtomicReference<>();
final AtomicReference<ListenableFuture<SegmentIdWithShardSpec>> futureReference = new AtomicReference<>();

// Possible race condition:
// t1 -> new batch is added to queue or batch already exists in queue
Expand Down Expand Up @@ -644,7 +644,7 @@ private class AllocateRequestBatch
* Map from allocate requests (represents a single SegmentAllocateAction)
* to the future of allocated segment id.
*/
private final Map<SegmentAllocateRequest, CompletableFuture<SegmentIdWithShardSpec>>
private final Map<SegmentAllocateRequest, SettableFuture<SegmentIdWithShardSpec>>
requestToFuture = new HashMap<>();

AllocateRequestBatch(AllocateRequestKey key)
Expand Down Expand Up @@ -692,10 +692,10 @@ boolean isFull()
return size() >= MAX_BATCH_SIZE;
}

Future<SegmentIdWithShardSpec> add(SegmentAllocateRequest request)
ListenableFuture<SegmentIdWithShardSpec> add(SegmentAllocateRequest request)
{
log.debug("Adding request to batch [%s]: %s", key, request.getAction());
return requestToFuture.computeIfAbsent(request, req -> new CompletableFuture<>());
return requestToFuture.computeIfAbsent(request, req -> SettableFuture.create());
}

void transferRequestsFrom(AllocateRequestBatch batch)
Expand All @@ -718,7 +718,7 @@ void failPendingRequests(Throwable cause)
{
if (!requestToFuture.isEmpty()) {
log.warn("Failing [%d] requests in batch[%s], reason[%s].", size(), key, cause.getMessage());
requestToFuture.values().forEach(future -> future.completeExceptionally(cause));
requestToFuture.values().forEach(future -> future.setException(cause));
requestToFuture.keySet().forEach(
request -> emitTaskMetric("task/action/failed/count", 1L, request)
);
Expand All @@ -732,7 +732,7 @@ void completePendingRequestsWithNull()
return;
}

requestToFuture.values().forEach(future -> future.complete(null));
requestToFuture.values().forEach(future -> future.set(null));
requestToFuture.keySet().forEach(
request -> emitTaskMetric("task/action/failed/count", 1L, request)
);
Expand All @@ -746,7 +746,7 @@ void handleResult(SegmentAllocateResult result, SegmentAllocateRequest request)
if (result.isSuccess()) {
emitTaskMetric("task/action/success/count", 1L, request);
IndexTaskUtils.emitSegmentAllocateMetric(result.getSegmentId(), request.getTask(), emitter);
requestToFuture.remove(request).complete(result.getSegmentId());
requestToFuture.remove(request).set(result.getSegmentId());
} else if (request.canRetry()) {
log.debug(
"Allocation failed on attempt [%d] due to error[%s]. Can still retry action[%s].",
Expand All @@ -759,7 +759,7 @@ void handleResult(SegmentAllocateResult result, SegmentAllocateRequest request)
+ " Completing action[%s] with a null value.",
request.getAttempts(), result.getErrorMessage(), request.getAction()
);
requestToFuture.remove(request).complete(null);
requestToFuture.remove(request).set(null);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.core.type.TypeReference;
import com.google.common.util.concurrent.ListenableFuture;
import org.apache.druid.error.DruidException;
import org.apache.druid.error.InvalidInput;
import org.apache.druid.indexing.common.TaskLock;
Expand Down Expand Up @@ -219,12 +220,19 @@ public SegmentPublishResult perform(Task task, TaskActionToolbox toolbox)
}

IndexTaskUtils.emitSegmentPublishMetrics(retVal, task, toolbox);
return retVal;
}

if (toolbox.shouldFailSegmentPublishImmediately(retVal, task, supervisorId, startMetadata)) {
return SegmentPublishResult.fail(retVal.getErrorMsg());
} else {
return retVal;
}
@Override
public boolean canPerformAsync(Task task, TaskActionToolbox toolbox)
{
return supervisorId != null && startMetadata != null;
}

@Override
public ListenableFuture<SegmentPublishResult> performAsync(Task task, TaskActionToolbox toolbox)
{
return toolbox.publishSegmentsWhenReady(task, supervisorId, startMetadata, this::perform);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import com.fasterxml.jackson.core.type.TypeReference;
import com.google.common.base.Preconditions;
import com.google.common.collect.ImmutableSet;
import com.google.common.util.concurrent.ListenableFuture;
import org.apache.druid.common.config.Configs;
import org.apache.druid.indexing.common.LockGranularity;
import org.apache.druid.indexing.common.TaskLock;
Expand Down Expand Up @@ -259,12 +260,19 @@ public SegmentPublishResult perform(Task task, TaskActionToolbox toolbox)
}

IndexTaskUtils.emitSegmentPublishMetrics(retVal, task, toolbox);
return retVal;
}

if (toolbox.shouldFailSegmentPublishImmediately(retVal, task, supervisorId, startMetadata)) {
return SegmentPublishResult.fail(retVal.getErrorMsg());
} else {
return retVal;
}
@Override
public boolean canPerformAsync(Task task, TaskActionToolbox toolbox)
{
return supervisorId != null && startMetadata != null;
}

@Override
public ListenableFuture<SegmentPublishResult> performAsync(Task task, TaskActionToolbox toolbox)
{
return toolbox.publishSegmentsWhenReady(task, supervisorId, startMetadata, this::perform);
}

private void checkWithSegmentLock()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,10 +22,12 @@
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import com.fasterxml.jackson.core.type.TypeReference;
import com.google.common.util.concurrent.ListenableFuture;
import org.apache.druid.indexing.common.task.Task;

import java.util.concurrent.Future;

/**
* An action performed on behalf of a Task by the Overlord.
*/
@JsonTypeInfo(use = JsonTypeInfo.Id.NAME, property = TaskAction.TYPE_FIELD)
@JsonSubTypes(value = {
@JsonSubTypes.Type(name = "lockAcquire", value = TimeChunkLockAcquireAction.class),
Expand Down Expand Up @@ -66,7 +68,7 @@ default boolean canPerformAsync(Task task, TaskActionToolbox toolbox)
return false;
}

default Future<RetType> performAsync(Task task, TaskActionToolbox toolbox)
default ListenableFuture<RetType> performAsync(Task task, TaskActionToolbox toolbox)
{
throw new UnsupportedOperationException();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,25 @@

package org.apache.druid.indexing.common.actions;

import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;

import java.io.IOException;

/**
* Client to submit a {@link TaskAction} on behalf of a Task to the Overlord.
*/
public interface TaskActionClient
{
<RetType> RetType submit(TaskAction<RetType> taskAction) throws IOException;

default <RetType> ListenableFuture<RetType> submitAsync(TaskAction<RetType> taskAction)
{
try {
return Futures.immediateFuture(submit(taskAction));
}
catch (Exception e) {
return Futures.immediateFailedFuture(e);
}
}
}
Loading
Loading