PANKAJ KUMAR
03/13/2024, 12:20 PMfailed to change its status from [READING] to [PAUSED], aborting. But for the first spike, there was no error.
I checked the logs; usually, 1 checkpoint happens at a time, but here, 3 checkpoints happened. A1 triggered the first checkpoint and then the other 2 by A2. The difference between the sequence offset for the first checkpoint is ~11M, and for the second and third checkpoints, it is less than 60. These 3 checkpoints cause 3 segment allocate requests.
-----Performing action for task[A1]: CheckPointDataSourceMetadataAction{supervisorId='abc', taskGroupId='3', checkpointMetadata=KafkaDataSourceMetadata{SeekableStreamStartSequenceNumbers=SeekableStreamStartSequenceNumbers{stream='xyz', partitionSequenceNumberMap={3=222871437, 183=81557685}, exclusivePartitions=[]}}}
Checkpointing [KafkaDataSourceMetadata{SeekableStreamStartSequenceNumbers=SeekableStreamStartSequenceNumbers{stream='xyz', partitionSequenceNumberMap={3=22871437, 183=81557685}, exclusivePartitions=[]}}] for taskGroup [3]
Handled checkpoint notice, new checkpoint is [{3=235525861, 183=94516712}] for taskGroup [3]
----Performing action for task[A2]: CheckPointDataSourceMetadataAction{supervisorId='abc', taskGroupId='3', checkpointMetadata=KafkaDataSourceMetadata{SeekableStreamStartSequenceNumbers=SeekableStreamStartSequenceNumbers{stream='xyz', partitionSequenceNumberMap={3=235525861, 183=94516712}, exclusivePartitions=[]}}}
Checkpointing [KafkaDataSourceMetadata{SeekableStreamStartSequenceNumbers=SeekableStreamStartSequenceNumbers{stream='xyz', partitionSequenceNumberMap={3=235525861, 183=94516712}, exclusivePartitions=[]}}] for taskGroup [3]
Handled checkpoint notice, new checkpoint is [{3=235525927, 183=94516992}] for taskGroup [3]
----Performing action for task[A2]: CheckPointDataSourceMetadataAction{supervisorId='abc', taskGroupId='3', checkpointMetadata=KafkaDataSourceMetadata{SeekableStreamStartSequenceNumbers=SeekableStreamStartSequenceNumbers{stream='xyz', partitionSequenceNumberMap={3=235525927, 183=94516992}, exclusivePartitions=[]}}}
Checkpointing [KafkaDataSourceMetadata{SeekableStreamStartSequenceNumbers=SeekableStreamStartSequenceNumbers{stream='xyz', partitionSequenceNumberMap={3=235525927, 183=94516992}, exclusivePartitions=[]}}] for taskGroup [3]
Handled checkpoint notice, new checkpoint is [{3=235525972, 183=94517492}] for taskGroup [3]
I still do not understand why multiple checkpoints happened?PANKAJ KUMAR
03/13/2024, 12:20 PMPANKAJ KUMAR
03/13/2024, 12:29 PMJohn Kowtko
03/13/2024, 3:35 PMPANKAJ KUMAR
03/13/2024, 5:38 PMAbhishek Agarwal
03/14/2024, 5:16 AMPANKAJ KUMAR
03/19/2024, 4:57 AMAmatya Avadhanula
03/19/2024, 12:24 PMPANKAJ KUMAR
03/21/2024, 3:26 PMJohn Kowtko
03/21/2024, 3:36 PMPANKAJ KUMAR
03/21/2024, 3:38 PMJohn Kowtko
03/21/2024, 3:43 PMPANKAJ KUMAR
04/02/2024, 5:16 AMPANKAJ KUMAR
04/15/2024, 1:33 PMMaxTotalRows case is the one which triggers multiple checkpoint. As soon as it hits maxtotalRows it update the SequenceToCheckpoint, when sequenceToCheckpoint is not null we create a new checkPointMetadataSourceAction. And after checkpointing, the task again start ingesting, And for my case the task was ingesting from multiple partitions and as soon as it reaches endoffset of the current sequence for 1 partition(endoffset was set during the checkpoint request, And since we have set the ingestion replica to 2 the endoffset won’t match with the current task offset), it start ingest in the new sequence and then again update the SequenceToCheckpoint to new sequence(as the new sequence is not checkpointed and maxTotalRows is still above the limit). We only reduced the maxTotalRows variable when we do persist. And from the code it looks to me we have disabled incremental persist.
// do not allow incremental persists to happen until all the rows from this batch
// of rows are indexed
So basically this triggers multiple checkpoint request. And these multiple checkpoint request can cause ingestion lag in 2 ways.
1. If the multiple checkpoint request is only for 1-2 task, it throttle ingestion becuase of multiple persist.
2. after a point, when multiple tasks issues multiple checkpoint requests, the ingestion tasks start failing with failed to change its status from [READING] to [PAUSED], aborting.
Now i do not understand why we have disabled incremental persist? If we do not enable incremental persist then how we can handle this case?
And i was able to replicate this issue in one of our testing druid cluster. (I only see multiple checkpoint in case we hit maxTotalRows, For maxRowPerSegment it works as expected).
cc: @Abhishek Agarwal @Amatya Avadhanula @Xavier @kfarazPANKAJ KUMAR
04/16/2024, 8:32 AM