Details
-
Bug
-
Status: Resolved
-
P2
-
Resolution: Fixed
-
None
-
None
Description
The KinesisReader implementing KinesisIO reports backlog by implementing the
UnboundedSource.getTotalBacklogBytes()
method as opposed to the
UnboundedSource.getSplitBacklogBytes()
This value is supposed to represent the total backlog across all shards. This function is implemented by calling SimplifiedKinesisClient.getBacklogBytes with the watermark of the kinesis shards managed within the UnboundedReader instance. As this watermark may be further ahead than the watermark across all shards, this may miss backlog bytes.
An additional concern is that the watermark is calculated using a WatermarkPolicy, which means that the watermark may be inconsistent to the kinesis timestamp for querying backlog.
Attachments
Issue Links
- links to