Skip to content

feat: Block segment publish action for a streaming task until prior task groups are complete - #20505

Open
kfaraz wants to merge 11 commits into
apache:masterfrom
kfaraz:segment_publish_queue
Open

kfaraz wants to merge 11 commits into
apache:masterfrom
kfaraz:segment_publish_queue

Conversation

@kfaraz

@kfaraz kfaraz commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor

Description

Follow up to #19091 as considered in #19091 (comment).

The changes in #19091 have shown that tasks often need to wait for their priors to finish publishing as we often see logs like:

Task[xyz] needs to wait before publishing as other taskGroups[abc, def]
are currently publishing to partition[KafkaTopicPartition{partition=1, topic='null', multiTopicPartition=false}].

Before this patch, if publish from streaming task A is slow, a later generation task B keeps trying to publish repeatedly until eventually failing once the number of retries is exhausted.
It is difficult to fix the number of retries since publish times may vary across supervisors and over time.
It is also pointless to retry if we already know that a prior task is yet to finish publishing.

This patch tries to address these problems by blocking the task action of the later task B until A is complete.

Changes

  • Add TaskActionClient.submitAsync which returns a ListenableFuture
  • Update API /druid/indexer/v1/action in OverlordResource to be async by employing the new submitAsync() method
  • Update task actions segmentTransactionalAppend and segmentTransactionalReplace to implement performAsync()
  • Update LocalTaskActionClient.submitAsync() to use performAsync() when the task action supports it
  • The segment publish is now attempted only when the task is ready to publish
  • The API response is sent only once the publish attempt is done
  • For non-streaming tasks, publish is attempted right away
  • Add method isTaskReadyToPublish() in SeekableStreamSupervisor.
    • This method returns a ListenableFuture that completes when all prior task groups have completed

Other changes

  • Update SegmentAllocateAction and SegmentAllocationQueue to use ListenableFuture to align with the new signature of performAsync()

Follow up items

  • Emit metric from SeekableStreamSupervisor and include the waiting tasks in the supervisor report to allow debugging in case a task remains stuck waiting to publish
  • Add more tests for new methods
  • Tie task action timeout to supervisor completion timeout

This PR has:

  • been self-reviewed.
  • added documentation for new or modified features or behaviors.
  • a release note entry in the PR description.
  • added Javadocs for most classes and all non-trivial methods. Linked related entities via Javadoc links.
  • added or updated version, license, or notice information in licenses.yaml
  • added comments explaining the "why" and the intent of the code wherever would not be obvious for an unfamiliar reader.
  • added unit tests or modified existing tests to cover new code paths, ensuring the threshold for code coverage is met.
  • added integration tests.
  • been tested in a test Druid cluster.

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

🟡 Changes recommended

Address the async response lifecycle failures before merging: ordinary actions returning null fail before serialization, and completed response writes never close the async request. The readiness path also has a completion-registration race and can run database publication on the supervisor's single-threaded notice executor.

Reviewed all 17 of 17 changed files, including the surrounding task-action client, supervisor lifecycle, and metadata publication paths.

Validation: git diff --check 131989817271bc836e815db738cc2be33fba3538 HEAD passed. Static review only; no tests or builds were run.

Severity Findings
P0 0
P1 2
P2 2
P3 0
Total 4

This is an automated review by Codex GPT-5.6-Luna(max)

After addressing the findings or replying to the comments, you can request another review from me to trigger a new automated review.

catch (Exception e) {
log.noStackTrace().warn(e, "Request timed out or is closed already");
}
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Complete the async response on every path

Finding: doAction starts an AsyncContext, but this helper only sets the status and optionally writes a body; it never completes the context. Consequently, even a finished task action leaves the HTTP exchange open until the async timeout, and the no-leader path returns here before installing any timeout listener.

Suggestion: Complete the async context after the response is written, including when no response body is present.

// It would be great to verify that this worker is actually supposed to be running the task before
// actually doing the action. Some ideas for how that could be done would be using some sort of attempt_id
// or token that gets passed around.
completeRequest(asyncContext, Status.OK.getStatusCode(), Map.of("result", result));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Preserve null task-action results

Finding: Several existing actions, including LockReleaseAction and UpdateLocationAction, implement TaskAction<Void> and return null. Map.of rejects null values, so these successful actions throw from onSuccess before completeRequest can send the result, leaving the async request unanswered.

Suggestion: Build the response with a null-tolerant map so the existing result: null response contract is preserved.

for (TaskGroup group : blockingTaskGroups.keySet()) {
group.addCompletionListener(() -> {
latch.countDown();
blockingTaskGroups.remove(group);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Avoid mutating the map being iterated

Finding: The constructor iterates blockingTaskGroups.keySet() while each completion callback removes from that same constructor argument. If a prior group completes during registration, addCompletionListener can run the callback inline (or concurrently), modifying the HashMap during iteration and throwing ConcurrentModificationException before the readiness future is returned.

Suggestion: Register listeners from an immutable snapshot and update a separate thread-safe tracking structure from the callbacks.

return SegmentPublishResult.retryableFailure("Task is not ready to publish yet");
}
},
MoreExecutors.directExecutor()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Keep publish work off the supervisor notice thread

Finding: The readiness future is completed by TaskGroup.onCompleted during checkPendingCompletionTasks, which runs on this supervisor's single-threaded notice executor. With a direct executor, this transform calls perform inline; that method enters the task-lock critical section and performs metadata-store publication, so a slow publish blocks this supervisor from processing subsequent status checks and notices. The same executor choice is present in the insert action.

Suggestion: Dispatch the publish operation to a dedicated executor when readiness completes, leaving the supervisor notice thread free to process its queue.

@Fly-Style
Fly-Style self-requested a review October 8, 2026 15:01
@kfaraz
kfaraz requested a review from FrankChen021 October 9, 2026 08:13

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

🟡 Changes recommended

The async publish path still has failure and timeout handling gaps that can leave action requests hanging or fail a streaming task before the prior task group finishes.

Reviewed 18 of 18 changed files, including the task-action client and HTTP response lifecycle, supervisor readiness/listener state, metadata publication path, and affected tests.

Validation: git diff --check 9a792495bd7e5ecdf6e8621b71b5983d072c3081 HEAD passed. Static review only; no tests or builds were run.

Severity Findings
P0 0
P1 2
P2 1
P3 0
Total 3

This is an automated review by Codex GPT-5.6 Luna(Max)

After addressing the findings or replying to the comments, you can request another review from me to trigger a new automated review.

MoreExecutors.directExecutor()
);
} else {
return Futures.immediateFuture(submit(taskAction));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Return synchronous action failures as failed futures

Finding: The fallback evaluates submit(taskAction) before immediateFuture can be created, and the async branch can likewise throw while evaluating canPerformAsync or performAsync. Those exceptions escape submitAsync instead of reaching the future callback. OverlordResource has already called startAsync() and invokes submitAsync before registering its timeout listener, so an ordinary action failure can leave the HTTP request without a response, timeout, or completion.

Suggestion: Convert synchronous failures from the async decision and both execution paths into failed futures, or guard the call in the resource so every failure completes the servlet context through the existing error callback.

public Response doAction(final TaskActionHolder holder)
public void doAction(
final TaskActionHolder holder,
@Nullable @PathParam("timeout") final Long timeoutMillis,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Make the async action timeout cover publish waits

Finding: The /action route has no {timeout} path segment, so @PathParam("timeout") is never populated and setTimeout always uses five minutes. The supervisor's default completionTimeout is 30 minutes; a legitimate pending task group that takes longer than five minutes causes the timeout listener to cancel the action and complete the request before the readiness future finishes. The worker then sees an I/O or empty-response failure, and the publish retry path does not retry that timeout, so the new blocking behavior can fail a task instead of waiting for the prior group.

Suggestion: Bind the timeout to a real request parameter or route and propagate a deadline that covers the configured supervisor completion wait, with an explicit retry-safe timeout response.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Typo: this should be query param.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Follow-up assessment

Agreed, @QueryParam("timeout") is the right binding. That fixes the annotation typo, but the timeout concern also needs the caller/default addressed: RemoteTaskActionClient posts to /druid/indexer/v1/action without a timeout query parameter, so changing the annotation alone still leaves worker requests capped at five minutes. A prior task group may legitimately remain pending under the supervisor's default 30-minute completionTimeout; the listener cancels the action instead of returning a retryable SegmentPublishResult.

Suggestion: Have worker requests use a deadline that covers the configured supervisor completion wait (or make the server derive it), and return an explicit response that the publish caller can safely retry when that deadline expires. The existing P1 remains actionable at this head; no additional finding is raised by this follow-up.

Reviewed 7 repository source files for this thread: 3 of the PR's 18 changed files and 4 supporting files. This was a focused static follow-up, not a full PR re-review.


This is an automated review by Codex GPT-5.6-Luna(max)

After addressing the findings or replying to the comments, you can request another review from me to trigger a new automated review.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Removed the parameter since the client is never sending it. A future PR may add a new field inside the task action payload itself. Alternatively, it may even be linked up to the completion timeout.

For now, on timeout, we send back a gateway timeout response so that the client retries (via StandardRetryPolicy).

completeRequest(
asyncContext,
Status.OK.getStatusCode(),
result == null ? Map.of() : Map.of("result", result)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Preserve the result field for null task-action results

Finding: The previous synchronous endpoint always serialized a result property, including {"result":null} for Void actions. The async success path now serializes {} when the action returns null. RemoteTaskActionClient tolerates the missing property, but callers that rely on the stable response shape will observe an API regression.

Suggestion: Use a null-tolerant response representation that retains the existing result field for successful null-valued actions.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

updated.

readyToPublish -> {
if (Boolean.TRUE.equals(readyToPublish)) {
// Task is already unblocked for publish, retrying will not fix offset mismatch
return doNotRetryOffsetMismatchFailure(publishAction.apply(task, this));

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.

nit: doNotRetryOffsetMismatchFailure name reads a bit ... confusing as for future success action. It would be great to have other name here; I personally did not get the idea from the first time.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

updated - made the logic inline instead of separate method.

this.supervisorManager = supervisorManager;
this.jsonMapper = jsonMapper;
this.segmentAllocationQueue = segmentAllocationQueue;
this.actionExec = scheduledExecutorFactory.create(4, "TaskActionToolbox-%s");

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 think, '4' threads worth a comment about the estimated load on publish readiness check.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

added, thanks for the suggestion.


return !getSupervisorManager().isAnotherTaskGroupPublishingToPartitions(
// Try publishing later if the failure was due to offset mismatch
final ListenableFuture<Boolean> taskReadyToPublishFuture = supervisorManager.isTaskReadyToPublishSegments(

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.

isTaskReadyToPublishSegments deeper may throw ISE or NotFound, should it be also wrapped by try-catch?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Yes, I am updating LocalTaskActionClient to handle it.

@kfaraz
kfaraz requested a review from FrankChen021 October 9, 2026 13:19

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

🟡 Changes recommended

The current async action path still has two correctness and lifecycle problems before merge: the initial publish attempt can execute synchronously before the servlet timeout is installed, and timeout cancellation races the direct future callback, so the intended 504 can become a 500 or a duplicate completion.

Reviewed 18 of 18 changed files in the full merge-base diff, including the surrounding task-action, supervisor, metadata-transaction, HTTP, and test code. I also rechecked both prior inline review threads against the current head. This was a static review; no build or tests were run. Validation performed: git diff --check 9a792495bd7e5ecdf6e8621b71b5983d072c3081 HEAD passed, and the local read-only merge gate passed against master.

Severity Findings
P0 0
P1 2
P2 0
P3 0
Total 2

This is an automated review by Codex GPT-5.6 Luna(Max)

After addressing the findings or replying to the comments, you can request another review from me to trigger a new automated review.

return;
}

final ListenableFuture<?> future = taskActionClient.get().submitAsync(holder.getAction());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Initial publish is outside the async timeout

Finding: This invokes submitAsync before the async listener and 15-minute timeout are installed. For the streaming publish actions handled by this PR, submitAsync calls performAsync synchronously and publishSegmentsWhenReady performs the first metadata/lock publish attempt inline before it returns a future. If that database or lock operation stalls, the servlet thread remains blocked and the request has no timeout or cancellation path, defeating the new async endpoint and allowing stalled publishes to consume Overlord request threads indefinitely.

Suggestion: Dispatch the initial publish through an executor or otherwise obtain the future without running the first attempt on the servlet thread, and install the timeout/cancellation handling so it covers that first attempt as well.

final ListenableFuture<?> future = taskActionClient.get().submitAsync(holder.getAction());
asyncContext.addListener(
ServletResourceUtils.createAsyncTimeoutListener(event -> {
future.cancel(true);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Timeout cancellation races the response callback

Finding: Canceling this future can synchronously invoke the onFailure callback because the callback is registered with directExecutor below; that callback treats cancellation as an internal-server-error response and completes the AsyncContext before this handler tries to write the intended 504. A concurrent success or failure can race the same timeout path as well, so one request can produce a 500/200 instead of the promised gateway timeout and then attempt a second completion on a closed response.

Suggestion: Serialize response completion with an atomic once guard, and make the timeout path claim the response or ignore the cancellation callback before canceling the underlying future.

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

🟡 Changes recommended

Two P1 async lifecycle races remain: the initial publish can run before timeout setup, and cancellation can trigger a 500 callback before the intended 504 timeout response.

Reviewed 18 of 18 changed files, including surrounding task-action, HTTP, supervisor, and metadata code. Static review only; no tests or builds were run.

Validation: git diff --check 9a792495bd7e5ecdf6e8621b71b5983d072c3081 HEAD passed.

Severity Findings
P0 0
P1 2
P2 0
P3 0
Total 2

This is an automated review by Codex GPT-5.6-Luna(max)

After addressing the findings or replying to the comments, you can request another review from me to trigger a new automated review.

return;
}

final ListenableFuture<?> future = taskActionClient.get().submitAsync(holder.getAction());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Initial publish runs before timeout setup

Finding: submitAsync is invoked before the timeout listener and 15-minute timeout are registered. Local streaming publish actions execute their first metadata publish synchronously inside submitAsync, so a slow first attempt is governed by the container's default timeout and cannot be canceled through this listener.

Suggestion: Install the timeout and one-shot completion state before starting the action, or move the initial publish attempt off the request thread.

asyncContext.addListener(
ServletResourceUtils.createAsyncTimeoutListener(event -> {
future.cancel(true);
completeAsyncRequest(event.getAsyncContext(), HttpResponseStatus.GATEWAY_TIMEOUT.code(), null);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P1] Cancellation can overwrite the timeout response

Finding: The timeout handler cancels the future before completing the 504 response. Because the callback uses directExecutor, cancel(true) can synchronously invoke onFailure with CancellationException; that branch treats cancellation as a generic 500 and completes the context before the timeout handler writes 504. A near-simultaneous completion can also cause duplicate response completion.

Suggestion: Use an atomic response-ownership guard, treat cancellation as the timeout path, and prevent late callbacks from writing after the 504 response.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants