Hi guys, we have recently enabled ingestion replic...
# general
p
Hi guys, we have recently enabled ingestion replica in our druid cluster. We saw an issue where the ingestion lag went high. I am trying to debug that, so I need some help. Here’s the scenario: The ingestion lag went high for both the replicas (A1, A2); there was no task failure. The ingestion lag started increasing after it finished the checkpointing. One difference is three checkpoints happened for the same set of partitions back-to-back. After all three checkpoints done, lag starts increasing on both A1, A2 and then after some time the segment transactional insert action happened and the A1, A2 start ingesting again and the lag came down on its own. But again, after a few cycles, the multiple checkpoints happened, and the lag again went high for the same partitions. Then, after task rollover, the lag started increasing on other partitions, and eventually, we restarted the overlord, which fixed this. For the last two ingestion lag spikes, we also saw an exception:
failed 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.
Copy code
-----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?
cc: @Abhishek Agarwal @Amatya Avadhanula @kfaraz @Xavier
On indexer logs: There is a lot of persist going on for the same task and also for other tasks running on that indexer, and I can see around the same time that there is a spike in persist time and persist backpressure. But I am not getting why persist impacted the ingestion? The ingestion was stopped for a few minutes and started ingesting again after the segmentTransactional insert action.
j
Hi Pankaj, Did you get an "ingestion throttled" statements in the task log? If so, then that means the consumer thread already filled up the next row buffer before the persist thread was able to finish persisting the prior row buffer. With maxPendingPersists set to 0 (the default), only one persist will be allowed to run concurrently, and if the consumer is too far ahead and fills up a second row buffer, then it will get throttled. You can try increasing maxPendingPersists from 0 to 1 or a higher number, but keep in mind you might just be prolonging the issue. If you are seeing the throttling, you can try increasing or decreasing the row buffer size (using maxRowsInMemory or intermediatePersistPeriod) to see if it changes the dynamic between the consumer and persist threads.
p
I don't see "ingestion throttled" in logs. But 'maxPendingPersists' is set to 0. I will try to increase this.
a
There is https://github.com/apache/druid/pull/13982 in 29 that could help too
p
Doubt: Is the persist triggered by the checkpointing? Or it is independent of checkpointing?
a
A persist is most commonly triggered when maxRowsInMemory for a task are hit. Checkpointing is an independent operation which is closely related to publishing of segments
🆗 1
p
Is this maxPendingPersist is per indexer or per ingestion task?
j
Per ingestion task
p
okay, So the spike on PersistTime and PersistBackpressure is because of more than 1 persist for this ingestion task or vice versa?
j
I haven't used persists/time before ... if this is like "total GC time" then I would think what you are suggesting, that it is the total of all persists happening, and if more than one persist is running in parallel then the persist time would be higher.
p
Hey, we again had this issue. And exactly same thing happened. A1 and A2 task checkpoints 3 time(first A1 checkpoint and then A2 and then again A2), And just after that multiple persist on indexer and ingestion got throttled. I looked into the codes to identify what causing the multiple checkpointing request. There are two cases in which ingestion task triggers checkpoint: 1. When (System.currentTimeMillis() nextCheckpointTime)>: We update the nextCheckpointTime everytime we set endoffset, so it will stop the other tasks from the taskgroup to trigger the checkpoint. 2. when we hit maxRowLimit or TotalRow: Here we check the number of rows, if it is higher than the row limit we trigger checkpoint and it also triggers the persist. I feel we are seeing multiple checkpoint request becuase of this case, as we also see multiple Persist. So here we use SequenceToUse, Can there be a case where it’s not updated when the other task have done the checkpointing?
Hey, I have debugged this further, I added some extra log lines to understand which case exactly triggered multiple checkpoints. And found out that The
MaxTotalRows
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.
Copy code
// 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 @kfaraz
Created a issue with a flow chart to explain it better. Please take a look. https://github.com/apache/druid/issues/16293