Hey Druid folks, <@U0319HD8HEC>, <@U031A9D19NU>, <...
# general
p
Hey Druid folks, @Amatya Avadhanula, @Abhishek Agarwal, @kfaraz We encountered a bug recently and need your help in understanding the root-cause. We are running Apache Druid 25.0.0 . We run our application with 2 ingestion replicas and we ingest data from Apache Kafka. In the steady state, we expect both the replica tasks to try and publish the same sequences at nearly the same time. However, we had two replica tasks, say D1 and D2. Now, the D2 task was publishing segments corresponding to sequences 0, 1, 2, 3 . However, task D1 directly tried to publish the segment corresponding to sequence 4. Prior to this, we had seen multiple checkpoints running for task D1, one for every sequence from 0 to 6. Now, we killed task D2 manually from the UI which had published segments corresponding to sequences 0, 1, 2, 3 to deep storage but not yet updated the metadata store.
Copy code
28 Mar 2024 @ 15:58:05.260 UTC Got shutdown request for task[index_kafka_<datasource>_cd0cec90105b458_<D2>]. Asking worker[<indexer-URL>] to kill it.

28 Mar 2024 @ 15:58:05.260 UTC Shutdown [index_kafka_<datasource>_cd0cec90105b458_<D2>] because: [An exception occurred while waiting for task [index_kafka_<datasource>_cd0cec90105b458_<D2>] to pause: [org.apache.druid.java.util.common.ISE: Task [index_kafka_<datasource>_cd0cec90105b458_<D2>] failed to change its status from [READING] to [PAUSED], aborting]]

28 Mar 2024 @ 15:58:05.261 UTC        Shutdown [index_kafka_<datasource>_cd0cec90105b458_<D2>] because: [shut down request via HTTP endpoint]
From the code, we see that it goes into this part:
Copy code
"// Publishing didn't affirmatively succeed. However, segments with our identifiers may still be active
              // now after all, for two possible reasons:
              //
              // 1) A replica may have beat us to publishing these segments. In this case we want to delete the
              //    segments we pushed (if they had unique paths) to avoid wasting space on deep storage.
              // 2) We may have actually succeeded, but not realized it due to missing the confirmation response
              //    from the overlord. In this case we do not want to delete the segments we pushed, since they are
              //    now live!"
and we see the log line:
Copy code
28 Mar 2024 @ 16:03:05.262 UTC.        Encountered exception in run() before persisting.
28 Mar 2024 @ 16:03:05.263 UTC	       Error while publishing segments for sequenceNumber[SequenceMetadata{sequenceId=0, sequenceName='index_kafka_<datasource>_cd0cec90105b458_0', assignments=[], startOffsets={226=1034753252873, 46=1680429116193}, exclusiveStartPartitions=[], endOffsets={226=1034767190850, 46=1680445277001}, sentinel=false, checkpointed=true}]
28 Mar 2024 @16:03:05.266 UTC.         Failed publish, not removing segments: [<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_104, <datasource>_2024-03-28T15:40:00.000Z_2024-03-28T15:50:00.000Z_2024-03-28T15:40:00.047Z_282]
28 Mar 2024 @ 16:08:39.422 UTC.        Found existing pending segment [<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_104] for sequence[index_kafka_<datasource>_cd0cec90105b458_0] (previous = [null]) in DB
Now, the successor tasks which are created for these tasks, say E1 and E2 (replicas of each other), are directly trying to publish the segments corresponding to sequence 5. Thus, they fail because the offset in the metadata store is still at startOffset of sequence_0 . So, now we are in a state, where, after every
intermediateHandoff
time, whenever the task tries to publish segments, it fails because the offsets do not match with the value stored in the metadata store. However, the tasks are able to get the sequence-segment mapping and are committing metadata corresponding to it. We also see
SegmentAllocateAction
run for all the sequences in E1 task. Sequence_0 corresponds to the segment
telemetry_metrics_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_104
. This info becomes available to the successor tasks, as can be seen here:
Copy code
28 Mar 2024 @ 16:35:42.453 UTC.  Committing metadata[AppenderatorDriverMetadata{segments={index_kafka_<datasource>_cd0cec90105b458_0=[SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:40:00.000Z_2024-03-28T15:50:00.000Z_2024-03-28T15:40:00.047Z_282, state=APPENDING}, SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_104, state=APPENDING}], index_kafka_<datasource>_cd0cec90105b458_1=[SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_135, state=APPENDING}], index_kafka_<datasource>_cd0cec90105b458_2=[SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_227, state=APPENDING}], index_kafka_<datasource>_cd0cec90105b458_3=[SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_231, state=APPENDING}], index_kafka_<datasource>_cd0cec90105b458_4=[SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_233, state=APPENDING}], index_kafka_<datasource>_cd0cec90105b458_5=[SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_254, state=APPENDING}], index_kafka_<datasource>_cd0cec90105b458_6=[SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_325, state=APPENDING}], index_kafka_<datasource>_cd0cec90105b458_7=[SegmentWithState{segmentIdentifier=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_329, state=APPENDING}]}, lastSegmentIds={index_kafka_<datasource>_cd0cec90105b458_0=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_104, index_kafka_<datasource>_cd0cec90105b458_1=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_135, index_kafka_<datasource>_cd0cec90105b458_2=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_227, index_kafka_<datasource>_cd0cec90105b458_3=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_231, index_kafka_<datasource>_cd0cec90105b458_4=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_233, index_kafka_<datasource>_cd0cec90105b458_5=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_254, index_kafka_<datasource>_cd0cec90105b458_6=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_325, index_kafka_<datasource>_cd0cec90105b458_7=<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_329}, callerMetadata={nextPartitions=SeekableStreamEndSequenceNumbers{stream='<streamName>', partitionSequenceNumberMap={226=1034768328624, 46=1680444294138}}}}] for sinks[<datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_227:1, <datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_104:64, <datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_325:1, <datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_135:8, <datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_254:2, <datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_233:1, <datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_231:1, <datasource>_2024-03-28T15:40:00.000Z_2024-03-28T15:50:00.000Z_2024-03-28T15:40:00.047Z_282:58, <datasource>_2024-03-28T15:50:00.000Z_2024-03-28T16:00:00.000Z_2024-03-28T15:50:00.039Z_329:1].
Thus, the sequenceMap had been updated by D2 task while the offsetMap had not been updated. We tried to restart overlord, but that did not help us mitigate this issue. We ultimately come out of this issue by setting the replicas to 1 and when the taskGroup (denoted by its prefix
cd0cec90105b458
) changes. The final state we arrive at is the same as the one described in this issue: https://github.com/apache/druid/issues/8139#issuecomment-801706838
Copy code
the system gets into a bad state where the supervisor is firing up new tasks based on the offsets that the previous tasks should've worked their way up to, but the previous tasks are not around to complete their job.
We have these queries: 1. How is this sequenceMap passed across tasks and how do the E1 task become aware that they should start publishing directly from sequence 5 instead of sequence 0 ? My expectation was that the follow-up tasks would try to publish all the sequences starting from sequence_0 as its offset matches with the value stored in the MDS. 2. Why is D1 directly publishing for sequence 4? It should also have published for sequences 0, 1, 2, 3. 3. What caused the sequenceMap and the offsetMap to go out of sync? Is the process of publishing segments and updating Metadata store with sequences and offsets not atomic? 4. Why do we have so many sequences created for such a small window? Even the offset difference is small for the later sequences. Who creates a new sequence? When is it triggered? 5. If we kill a task manually from the UI, is the atomicity of the transaction respected? cc: @Xavier, @PANKAJ KUMAR