{"record":{"id":"a9e83d9281529012","repo":"apache/druid","slug":"mismatched-kafka-and-task-partitions-missing-task","errorCode":null,"errorMessage":"Mismatched kafka and task partitions: Missing Task Partitions %s, Missing Kafka Partitions %s","messagePattern":"Mismatched kafka and task partitions: Missing Task Partitions (.+?), Missing Kafka Partitions (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisor.java","lineNumber":285,"sourceCode":"    return taskList;\n  }\n\n  @Override\n  protected Map<KafkaTopicPartition, Long> getPartitionRecordLag()\n  {\n    OffsetSnapshot<KafkaTopicPartition, Long> offsetSnapshot = offsetSnapshotRef.get();\n    Map<KafkaTopicPartition, Long> latestSequencesFromStream = offsetSnapshot.getLatestOffsetsFromStream();\n    Map<KafkaTopicPartition, Long> highestIngestedOffsets = offsetSnapshot.getHighestIngestedOffsets();\n\n    if (latestSequencesFromStream.isEmpty()) {\n      return null;\n    }\n\n    Set<KafkaTopicPartition> kafkaPartitions = latestSequencesFromStream.keySet();\n    Set<KafkaTopicPartition> taskPartitions = highestIngestedOffsets.keySet();\n    if (!kafkaPartitions.equals(taskPartitions)) {\n      try {\n        log.warn(\"Mismatched kafka and task partitions: Missing Task Partitions %s, Missing Kafka Partitions %s\",\n                sortingMapper.writeValueAsString(Sets.difference(kafkaPartitions, taskPartitions)),\n                 sortingMapper.writeValueAsString(Sets.difference(taskPartitions, kafkaPartitions)));\n      }\n      catch (JsonProcessingException e) {\n        throw DruidException.defensive(\"Failed to serialize KafkaTopicPartition when getting partition record lag: %s\",\n                                       e.getMessage());\n      }\n    }\n\n    return getRecordLagPerPartitionInLatestSequences(offsetSnapshot);\n  }\n\n  @Nullable\n  @Override\n  protected Map<KafkaTopicPartition, Long> getPartitionTimeLag()\n  {\n    return partitionToTimeLag;\n  }","sourceCodeStart":267,"sourceCodeEnd":303,"githubUrl":"https://github.com/apache/druid/blob/9b90983fd291f26935af934383ce360473179e4d/extensions-core/kafka-indexing-service/src/main/java/org/apache/druid/indexing/kafka/supervisor/KafkaSupervisor.java#L267-L303","documentation":"KafkaSupervisor.getPartitionRecordLag() logs a warning when the set of partitions the Kafka stream reports differs from the set of partitions currently assigned to tasks. It compares latestSequencesFromStream keys against highestIngestedOffsets keys and dumps the missing sets. Lag values cannot be computed per-partition when the sets disagree.","triggerScenarios":"Calling getPartitionRecordLag() (via partitionRecordLag) when Kafka has partitions with no active task, or tasks track partitions that no longer exist in the topic (partition count changed or task lag).","commonSituations":"Topic partition count increased but tasks not yet reassigned; partition count decreased (unsupported shrink); supervisor still starting up or tasks in unstable state; monitoring scrape during reassignment.","solutions":["Wait for the supervisor to converge; this is often transient during partition reassignment.","If partitions were added, verify taskCount/multiStageEngagement and let the supervisor spawn tasks for new partitions.","Never decrease a Kafka topic's partition count; if it happened, restore it or rebuild the datasource.","Check supervisor task state via /druid/indexer/supervisor status endpoints for stuck assignments."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"Set<Long> kafka = latestSequencesFromStream.keySet();\nSet<Long> task = highestIngestedOffsets.keySet();\nif (!kafka.equals(task)) { /* skip lag computation or alert on persistent mismatch */ }","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Monitor supervisor convergence after partition changes","Never shrink Kafka topic partition counts","Delay lag-based autoscaling until task assignments are stable"],"tags":["kafka","partitions","supervisor","monitoring"],"backgroundTag":"unexpected-response-shape","analyzedSha":"9b90983fd291f26935af934383ce360473179e4d","analyzedAt":"2026-09-07T13:32:30.957Z","contentChangedAt":"2026-09-07T13:32:30.957Z","schemaVersion":2},"datasetVersion":"2026-09-17T15:17:12.973Z"}