apache / apache/druid

Checkpointing failure after taskDuration in Kafka/Kinesis indexing service

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

Description

### Affected Version

All versions since incremental handoff was introduced

### Description

Checkpointing can be initiated by both the supervisor and tasks. Tasks can initiate checkpointing whenever it wants to publish segments. The supervisor initiates checkpointing when the task run time has reached to `taskDuration`. When the supervisor initiates checkpointing, the task changes its status to `publishing` and will stop once it publishes all segments.

The supervisor calls `checkTaskDuration()` to start checkpointing (https://github.com/apache/incubator-druid/blob/master/indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java#L1912-L1969).

Here is some code snippet.

```java
private void checkTaskDuration() throws ExecutionException, InterruptedException, TimeoutException
{
final List>> futures = new ArrayList<>();
...
for (Entry entry : activelyReadingTaskGroups.entrySet()) {
...
if (earliestTaskStart.plus(ioConfig.getTaskDuration()).isBeforeNow()) {
log.info("Task group [%d] has run for [%s]", groupId, ioConfig.getTaskDuration());
futureGroupIds.add(groupId);
futures.add(checkpointTaskGroup(group, true));
}
}

List> results = Futures.successfulAsList(futures)
.get(futureTimeoutInSeconds, TimeUnit.SECONDS);
for (int j = 0; j < results.size(); j++) {
...
activelyReadingTaskGroups.remove(groupId);
}
}
```

The issue is `checkpointTaskGroup(group, true)` can be called more than one time for the same taskGroup if `Futures.successfulAsList(futures).get(futureTimeoutInSeconds, TimeUnit.SECONDS)` fails because of timeout. If it is timed out, some future might fail, but others might succeed. However, `activelyReadingTaskGroups` is updated only when `futures.get()` is returned successfully. As a result, when `checkTaskDuration` is called in the next runNotice, it can start duplicate checkpointing for some taskGroups because they are still in `activelyReadingTaskGroups` which results in failing all tasks in those taskGroups because the previous checkpointing succeeded and they are now in `publishing` status.

Contributor guide

Open the contributing guide

Research direction

Start in indexing-service/src/main/java/org/apache/druid/indexing/seekablestream/supervisor/SeekableStreamSupervisor.java, especially checkTaskDuration() and checkpointTaskGroup(). Examine the timeout path around Futures.successfulAsList(...).get(...) and how activelyReadingTaskGroups is updated. Done means a timed-out check cannot cause duplicate checkpointing for task groups whose earlier checkpoint succeeded.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, kafka
Domain
backend, stream-processing
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.