Slackbot
02/22/2024, 4:47 PMAmatya Avadhanula
02/23/2024, 1:58 AMParth Agrawal
02/23/2024, 3:32 AMA2 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.Amatya Avadhanula
02/23/2024, 3:35 AMSELECT COUNT(*) FROM sys."segments"
Total segments per interval:
SELECT COUNT(*) as num_segments, "start", "end"
FROM sys."segments"
GROUP BY 2, 3
ORDER BY 1 DESC, 2 DESCParth Agrawal
03/04/2024, 3:27 PMshutting 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:
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?Parth Agrawal
03/04/2024, 3:28 PMTotal segments in cluster:Around 150k segments total
Abhishek Agarwal
03/04/2024, 4:45 PMParth Agrawal
03/05/2024, 4:55 AMParth Agrawal
03/12/2024, 3:36 AMAbhishek Agarwal
03/12/2024, 2:01 PMAbhishek Agarwal
03/12/2024, 2:02 PMParth Agrawal
03/20/2024, 3:54 PMAbhishek Agarwal
03/21/2024, 11:04 AMParth Agrawal
04/22/2024, 4:26 AMAmatya Avadhanula
04/22/2024, 10:28 AMParth Agrawal
05/06/2024, 4:36 AMParth Agrawal
05/06/2024, 4:37 AMWhen 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.
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.Amatya Avadhanula
05/06/2024, 4:38 AMAmatya Avadhanula
05/06/2024, 4:38 AMParth Agrawal
05/06/2024, 4:39 AMParth Agrawal
05/06/2024, 4:40 AMParth Agrawal
05/06/2024, 4:42 AMParth Agrawal
05/06/2024, 4:42 AMAmatya Avadhanula
05/06/2024, 4:44 AMParth Agrawal
05/06/2024, 4:45 AMParth Agrawal
05/07/2024, 4:49 AM1. Trigger shutdownThis refers to killing A2' task from UI ?
Amatya Avadhanula
05/07/2024, 4:50 AMAmatya Avadhanula
05/07/2024, 4:51 AMParth Agrawal
05/07/2024, 4:52 AMParth Agrawal
05/07/2024, 5:04 AMAmatya Avadhanula
05/07/2024, 5:15 AMParth Agrawal
05/07/2024, 5:18 AMParth Agrawal
05/07/2024, 5:19 AMAmatya Avadhanula
05/07/2024, 5:19 AMParth Agrawal
05/07/2024, 6:11 AMAmatya Avadhanula
05/07/2024, 6:11 AMParth Agrawal
05/07/2024, 8:36 AMParth Agrawal
05/14/2024, 4:36 AMentireTaskGroupFailed works.
As per my understanding of the 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.
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?