apache / apache/druid

KafkaIndexingTask pending forever after restart in remote mode

Open
#8,580 4 comments 0 reactions 0 assignees View on GitHub
Area - Ingestion Bug
Dominant language
Java
Stars
14.1k
Forks
3.8k
Avg merge
2d 58m
Merged PRs (30d)
233

Description

### Affected Version

All

### Description

#### Cluster size

2Master Node,2 Query Node,4Data Node
In fact,if you want to reproduce the problem , you are advised to use only 1 overlord and 1 middlemanager.

#### Configurations in use

druid.indexer.runner.type = remote

default configurations

### Reason

#### Code 1:
org.apache.druid.indexing.overlord.TaskQueue#manage
```
while (active) {
for (final Task task : ImmutableList.copyOf(tasks)) {
if (!taskFutures.containsKey(task.getId())) {
// doSomething...
taskFutures.put(task.getId(), attachCallbacks(task, runnerTaskFuture));
}
}

managementMayBeNecessary.awaitNanos(60s);
}
```
If a task both in tasks and taskFutures , it will be considered to be monitored with callbacks; when it status changed,it will notify mysql to sync status from zk.

**But what if it status changed, and callbacks did not work ?**

#### Code 2:

org.apache.druid.indexing.overlord.RemoteTaskRunner#addWorker

```
zkWorker.addListener(

new PathChildrenCacheListener()
{

public void childEvent(CuratorFramework client, PathChildrenCacheEvent event)
{
...

case CHILD_UPDATED:

...

if ((tmp = runningTasks.get(taskId)) != null) {
taskRunnerWorkItem = tmp;
} else {
final RemoteTaskRunnerWorkItem newTaskRunnerWorkItem = new RemoteTaskRunnerWorkItem(
taskId,
announcement.getTaskType(),
zkWorker.getWorker(),
TaskLocation.unknown(),
announcement.getTaskDataSource()
);
final RemoteTaskRunnerWorkItem existingItem = runningTasks.putIfAbsent(
taskId,
newTaskRunnerWorkItem
);
if (existingItem == null) {
log.warn(
"Worker[%s] announced a status for a task I didn't know about, adding to runningTasks: %s",
zkWorker.getWorker().getHost(),
taskId
);
taskRunnerWorkItem = newTaskRunnerWorkItem;
} else {
taskRunnerWorkItem = existingItem;
}
}
...
if (announcement.getTaskStatus().isComplete()) {
taskComplete(taskRunnerWorkItem, zkWorker, announcement.getTaskStatus());
runPendingTasks();
}
}
}
)
```
if the task in zk not exisit in runningTasks, it will new a taskRunnerWorkItem without attch callbacks on it .And if the task status change to failed,it will execute the new callback . So the task
miss to sync zk status to mysql,but it still stay in tasks and taskFutures.

it will be pending forever.

### Steps to reproduce the problem

1. start overlord ,start coordinator ,start middlemanager ,start historical ,start broker .
2. when there is a task running in middlemanager ,stop middlemanager ,and stop overlord after middlemanager.
3. a few minutes later,start overlord.
4. tail -f overlord's log , when you see "Beginning management".it means start code 1, 30s or 90s after it appeared, start middlemanager.
5. after middlemanager started , there would be a task in pending status.

I hava already reproduce the problem.
After I dump the overlord, the task exist is really both in tasks and taskFutures list.
The pending task will be clear after killed and reset the supervisor.

Contributor guide

Open the contributing guide

Research direction

Start with TaskQueue#manage and RemoteTaskRunner#addWorker, focusing on how tasks and taskFutures are rebuilt after an overlord and middlemanager restart. Reproduce the documented stop/start sequence with one overlord and one middlemanager, then verify that the task status is synchronized from ZooKeeper and does not remain pending indefinitely.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
backend, distributed-systems
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.