apache / apache/druid

Kinesis supervisor is showing unhealthy and task are not running when one or more partition is empty

Open
#9,763 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
0.18.0

### Description

There is no problem when index task is running and polling from kafka/kinesis stream with one or more empty shards (as tested in KafkaIndexTaskTest.java and KinesisIndexTaskTest.java). The problem for Kinesis described in the tittle is when we try to get the sequence number in SeekableStreamSupervisor#getOffsetFromStorageForPartition and Kinesis has one or more empty shard (as tested in KinesisRecordSupplierTest.java and SeekableStreamSupervisorStateTest.java). More specifically, this happens for the following conditions:

- we don't have a startingOffset (first run or we had some previous failures and reset the sequences) and don't have offset in metadata store so we retrieve the latest or earliest Kinesis sequence

- we don't have a startingOffset (first run or we had some previous failures and reset the sequences) and we have offset in metadata store but skipSequenceNumberAvailabilityCheck=False

Currently for Kinesis, in SeekableStreamSupervisor#getOffsetFromStorageForPartition, after we use a ShardIterator to get some records, you get back a new iterator to continue reading where you left off. The thing is, it doesn't matter whether or not you've already reached the end of the stream, you'll still get back a valid ShardIterator. As long as the shard is open, any call to GetRecords with a valid (unexpired) ShardIterator will provide a valid non-null NextShardIterator. Hence, we keep getting new ShardIterator until we timeout and then throw an ISE exception which then resulted in Kinesis supervisor showing unhealthy.

We should determine when Kinesis shard is empty and not rely on timeout.

Contributor guide

Open the contributing guide

Research direction

Start with SeekableStreamSupervisor#getOffsetFromStorageForPartition and compare the empty-shard cases described in KinesisRecordSupplierTest.java and SeekableStreamSupervisorStateTest.java; KinesisIndexTaskTest.java provides the non-failing comparison. Reproduce the empty-shard behavior and verify that the supervisor determines the shard is empty without waiting for a timeout or becoming unhealthy.

Written by the indexing model from the issue text.

Assessment

Tech stack
aws, java
Domain
data-engineering, 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.