{"record":{"id":"86555c795668f4fc","repo":"apache/seatunnel","slug":"partition-of-topic-has-no-leader-skipping-d","errorCode":null,"errorMessage":"Partition {} of topic {} has no leader, skipping due to ignore_no_leader_partition=true.","messagePattern":"Partition (.+?) of topic (.+?) has no leader, skipping due to ignore_no_leader_partition=true\\.","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaSourceSplitEnumerator.java","lineNumber":395,"sourceCode":"                currentPathTopics.addAll(Arrays.asList(metadata.getTopic().split(\",\")));\n            }\n            currentPathTopics.forEach(topic -> topicMappingTablePathMap.put(topic, tablePath));\n            topics.addAll(currentPathTopics);\n        }\n        log.info(\"Discovered topics: {}\", topics);\n        Collection<TopicPartition> partitions =\n                adminClient.describeTopics(topics).allTopicNames().get().values().stream()\n                        .flatMap(\n                                t ->\n                                        t.partitions().stream()\n                                                .filter(\n                                                        partitionInfo -> {\n                                                            if (kafkaSourceConfig != null\n                                                                    && kafkaSourceConfig\n                                                                            .isIgnoreNoLeaderPartition()\n                                                                    && partitionInfo.leader()\n                                                                            == null) {\n                                                                log.warn(\n                                                                        \"Partition {} of topic {} has no leader, skipping due to ignore_no_leader_partition=true.\",\n                                                                        partitionInfo.partition(),\n                                                                        t.name());\n                                                                return false;\n                                                            }\n                                                            return true;\n                                                        })\n                                                .map(\n                                                        p ->\n                                                                new TopicPartition(\n                                                                        t.name(), p.partition())))\n                        .collect(Collectors.toSet());\n        Map<TopicPartition, Long> latestOffsets = listOffsets(partitions, OffsetSpec.latest());\n        return partitions.stream()\n                .map(\n                        partition -> {\n                            // Obtain the corresponding topic TablePath from kafka topic\n                            TablePath tablePath = topicMappingTablePathMap.get(partition.topic());","sourceCodeStart":377,"sourceCodeEnd":413,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaSourceSplitEnumerator.java#L377-L413","documentation":"During split enumeration, a partition of a topic has no leader replica (e.g. after broker failure before leader election completes). Because ignore_no_leader_partition=true is set, the enumerator skips the leaderless partition instead of failing discovery.","triggerScenarios":"getTopicInfo (via fetchPendingPartitionSplit) filters partitionInfo entries; kafkaSourceConfig.isIgnoreNoLeaderPartition() is true and partitionInfo.leader() == null, so the partition is filtered out of the split list.","commonSituations":"Broker crash with unclean leadership loss; topic replication factor issues after a broker decommission; Kafka cluster under maintenance/restart; newly created topic partitions without elected leaders.","solutions":["Restore the leader: restart/recover the failed broker or trigger preferred leader election","Fix topic health (increase replication factor, run kafka-leader-election.sh --election-type PREFERRED)","Disable ignore_no_leader_partition if you prefer failing fast over silent skips","Re-run or restart the job after the cluster recovers so skipped partitions are picked up"],"exampleFix":"// before\nKafkaSource {\n  ignore_no_leader_partition = true   # partitions silently skipped\n}\n// after\nKafkaSource {\n  ignore_no_leader_partition = false  # fail fast when leaders are missing\n}","handlingStrategy":"validation","validationCode":"// Verify topic leadership health before submitting the job\nAdminClient admin = AdminClient.create(props);\nMap<TopicPartition, Optional<LeaderAndIsr>> desc =\n    admin.describeTopics(Collections.singleton(topic)).allTopicNames().get();\n// fail submission if any partition lacks a leader","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Set adequate replication factor (>=3) so leaders can be re-elected","Monitor under-min-ISR / offline partitions in the cluster","Prefer failing fast (ignore_no_leader_partition=false) unless skips are expected","Run preferred leader election after broker maintenance"],"tags":["kafka","partition","discovery"],"backgroundTag":"no-leader-partition","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}