This message was deleted.
# general
s
This message was deleted.
a
For now, it would be helpful to see why the completion timeout elapses and try to mitigate it by increasing the configured value. I suspect it is the handoff taking a long time. Could you please share the total number of segments on the cluster and the max number of segments per interval?
p
Yes, I am debugging on that end. The task report file is <400 bytes. The problematic task
A2
never pushed any task report but when the main task runner thread was interrupted with exception
Encountered exception while waiting for task [A2] shutdown
, it was stuck at pushing the task reports to S3.
a
Please try these queries in the Druid console: Total segments in cluster:
Copy code
SELECT COUNT(*) FROM sys."segments"
Total segments per interval:
Copy code
SELECT COUNT(*) as num_segments, "start", "end"
FROM sys."segments"
GROUP BY 2, 3
ORDER BY 1 DESC, 2 DESC
p
I debugged the logs further and this is the flow of events: • Handoff started for A2 task at 15:23. The task was waiting for handoff to be complete. • At 152619 , A1 completed and gave the stop signal to A2. When A2 received the stop signal, it started
shutting down immediately
and dropped the segment which it was handling. The segment was dropped successfully. The task could not have been waiting for a segment to pick up the segment as the segment was already loaded to historicals by the replica task which had completed. • At 152623 , discoverTasks() function ran and put A2 in a new
Pending Completion
task group. • The task never completed shutting down and was stuck somewhere till the timeout elapsed. I see no logs for the task coming between 15:27 and 15:56 except stating its current offset. • At 15:56, we see this exception:
Copy code
java.util.concurrent.ExecutionException: java.lang.RuntimeException: java.lang.RuntimeException: <http://org.apache.druid.java.util.common.RE|org.apache.druid.java.util.common.RE>: Current thread is interrupted after [0] tries
	at java.base/java.util.concurrent.FutureTask.report(FutureTask.java:122)
	at java.base/java.util.concurrent.FutureTask.get(FutureTask.java:205)
	at org.apache.druid.indexing.overlord.ThreadingTaskRunner$2.call(ThreadingTaskRunner.java:323)
	at org.apache.druid.indexing.overlord.ThreadingTaskRunner$2.call(ThreadingTaskRunner.java:315)
	at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
	at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
	at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
	at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: java.lang.RuntimeException: java.lang.RuntimeException: <http://org.apache.druid.java.util.common.RE|org.apache.druid.java.util.common.RE>: Current thread is interrupted after [0] tries
	at org.apache.druid.indexing.overlord.ThreadingTaskRunner$1.call(ThreadingTaskRunner.java:232)
	at org.apache.druid.indexing.overlord.ThreadingTaskRunner$1.call(ThreadingTaskRunner.java:152)
	... 4 more
Caused by: java.lang.RuntimeException: <http://org.apache.druid.java.util.common.RE|org.apache.druid.java.util.common.RE>: Current thread is interrupted after [0] tries
	at org.apache.druid.storage.s3.S3TaskLogs.pushTaskFile(S3TaskLogs.java:156)
	at org.apache.druid.storage.s3.S3TaskLogs.pushTaskReports(S3TaskLogs.java:141)
	at org.apache.druid.indexing.overlord.ThreadingTaskRunner$1.call(ThreadingTaskRunner.java:223)
	... 5 more
Caused by: <http://org.apache.druid.java.util.common.RE|org.apache.druid.java.util.common.RE>: Current thread is interrupted after [0] tries
	at org.apache.druid.java.util.common.RetryUtils.retry(RetryUtils.java:148)
	at org.apache.druid.java.util.common.RetryUtils.retry(RetryUtils.java:81)
	at org.apache.druid.java.util.common.RetryUtils.retry(RetryUtils.java:163)
	at org.apache.druid.java.util.common.RetryUtils.retry(RetryUtils.java:153)
	at org.apache.druid.storage.s3.S3Utils.retryS3Operation(S3Utils.java:101)
	at org.apache.druid.storage.s3.S3TaskLogs.pushTaskFile(S3TaskLogs.java:147)
	... 7 more
I have gone through the code but could not pinpoint where the task thread was stuck or the exception was swallowed. Can you please take a look?
Total segments in cluster:
Around 150k segments total
a
you don't have the stack trace of A2 during 15:27 to 15:56 by any chance, do you?
p
Nope
Bumping up this thread @Amatya Avadhanula @Abhishek Agarwal
a
It will be hard to troubleshoot this without knowing what that task was stuck on.
which is something a thread dump can tell us.
p
Hey @Abhishek Agarwal, While we can’t troubleshoot why A2 was stuck in handoff, Can we please discuss ways to resolve this issue: https://github.com/apache/druid/issues/15944 ? Even if a task times out, it should not cause the activelyReading tasks to be killed if the timed out task’s replica has already succeeded.
a
@Parth Agrawal - do you have an approach in mind?
p
Regarding https://github.com/apache/druid/issues/15944 , We want to kill the activelyReading taskGroups only if the publishing has not happened for the offsets which A2 was supposed to ingest. For this, can we fetch currentOffset in metadataStore for that particular partition? If that is more than A2’s endOffset that implies that the ingestion took place successfully till A2’s end offset and beyond. Does this approach sound good to you? Will this require the indexer to talk to overlord to get the metadata store value or is that information present with the indexers ?
a
> However, A2 is unable to complete publishing and thus stop. Is it unable to push to deep storage or is it stuck waiting for handoff?
p
The task was stuck at pushing to deep storage
I thought of this solution. Can you please review? Fetch currentOffset from MDS
Copy code
When do we want the activelyReading taskGroups to be killed?

If the publishing has not happened for the offsets which A2 was supposed to ingest.

Fetch currentOffset in metadataStore for that particular partition. If that is more than A2’s endOffset that implies that the ingestion took place successfully till A2’s end offset and beyond.

If the currentOffset values are present in memory with the overlord, there is no need to fetch this data from the RDS. We can simply use the values present with overlord.
Copy code
Comparison: Before we kill actively reading tasks, we can compare pendingCompletionTask.endOffsets() and currentOffset for that particular partition. If currentOffset is higher, then do not delete activelyReading task groups.
a
In such a case your proposed solution might not work. The sequence of events is to: build segments, push to deep storage, commit to metadata store, wait for handoff.
While your solution is valid, and might be good to have, it might not help you in this particular case
p
In this case, the other replica task (A1) had published and completed successfully, which means the metadata store offset was already >= the value A2 was trying to commit.
That makes A2 irrelevant as there is no data loss being caused by A2's failure
If A1 had failed or was also stuck, the metadata store would not have been updated
in which case we actually want to kill the successor tasks to reduce the amount of ingestion lag seen.
a
Could you please try out this scenario locally and see if you observe lag?: 1. Start ingestion with two replicas A1, A2 with task duration 10 mins 2. After 9 mins kill replica A2 3. Wait for replica replacement A2' to spawn 4. Trigger shutdown 5. After A2' succeeds, kill A1
p
ok, trying
1. Trigger shutdown
This refers to killing A2' task from UI ?
a
No. I meant graceful shutdown. You can try suspending the supervisor and running it again
I will try this shortly as well, along with a few other scenarios
p
Okay, trying this out, but during the issue that I had faced, one of the replicas (A1) had completed successfully and updated the metadata store.
> 1. After A2' succeeds, kill A1 A1 succeeded almost at the same time as A2' . I had tried suspending the supervisor and running it again. Is there a way to trigger gracefulShutdown for only A2' ?
a
There isn't. Is it a low volume stream?
p
Yes, the testing env is.
The actual stream is high volume where the issue was observed.
👍 1
a
Ah, in that case both of them would complete at the same time. The scenario is likely to be reproducible only in a high volume env
p
Can we use this API: https://druid.apache.org/docs/latest/api-reference/tasks-api/#shut-down-a-task to trigger graceful shutdown of a specific task?
a
I think it kills the task
p
I don’t have a test setup in a high volume env. If my proposed solution is not disruptive, shall I go ahead and create a PR for the same?
I am a bit confused on how the
entireTaskGroupFailed
works. As per my understanding of the code,
Copy code
if (taskData.status.isFailure()) {
            stateManager.recordCompletedTaskState(TaskState.FAILED);
            iTask.remove(); // remove failed task
            if (group.tasks.isEmpty()) {
              // if all tasks in the group have failed, just nuke all task groups with this partition set and restart
              entireTaskGroupFailed = true;
              break;
            }
          }
If one of the tasks in the pending completion task group has failed and that was the only task in the pending completion task group, then we mark it as true and then kill all the successor tasks. In the issue I faced, A2 was the only task in pendingCompletionTaskGroup and thus, this variable would have been set to true.
Copy code
if (entireTaskGroupFailed) {
            log.warn("All tasks in group [%d] failed to publish, killing all tasks for these partitions", groupId);
          } else {
            log.makeAlert(
                "No task in [%s] for taskGroup [%d] succeeded before the completion timeout elapsed [%s]! "
                + "Check metrics and logs to see if the creation, publish or handoff"
                + " of any segment is taking longer than usual.",
                group.taskIds(),
                groupId,
                ioConfig.getCompletionTimeout()
            ).emit();
          }
However, I saw the log line corresponding to the else part in this code, which means that entireTaskGroupFailed was set to false. Am I missing something?