apache/seatunnel · warning

Partition of topic has no leader, skipping due to…

Error message

Partition {} of topic {} has no leader, skipping due to ignore_no_leader_partition=true.

What it means

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.

Solutions

  1. Restore the leader: restart/recover the failed broker or trigger preferred leader election
  2. Fix topic health (increase replication factor, run kafka-leader-election.sh --election-type PREFERRED)
  3. Disable ignore_no_leader_partition if you prefer failing fast over silent skips
  4. Re-run or restart the job after the cluster recovers so skipped partitions are picked up

Example fix

// before
KafkaSource {
  ignore_no_leader_partition = true   # partitions silently skipped
}
// after
KafkaSource {
  ignore_no_leader_partition = false  # fail fast when leaders are missing
}
Defensive patterns

Strategy: validation

Validate before calling

// Verify topic leadership health before submitting the job
AdminClient admin = AdminClient.create(props);
Map<TopicPartition, Optional<LeaderAndIsr>> desc =
    admin.describeTopics(Collections.singleton(topic)).allTopicNames().get();
// fail submission if any partition lacks a leader

Prevention

When it happens

Trigger: getTopicInfo (via fetchPendingPartitionSplit) filters partitionInfo entries; kafkaSourceConfig.isIgnoreNoLeaderPartition() is true and partitionInfo.leader() == null, so the partition is filtered out of the split list.

Common situations: 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.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/86555c795668f4fc. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-kafka/src/main/java/org/apache/seatunnel/connectors/seatunnel/kafka/source/KafkaSourceSplitEnumerator.java:395

                currentPathTopics.addAll(Arrays.asList(metadata.getTopic().split(",")));
            }
            currentPathTopics.forEach(topic -> topicMappingTablePathMap.put(topic, tablePath));
            topics.addAll(currentPathTopics);
        }
        log.info("Discovered topics: {}", topics);
        Collection<TopicPartition> partitions =
                adminClient.describeTopics(topics).allTopicNames().get().values().stream()
                        .flatMap(
                                t ->
                                        t.partitions().stream()
                                                .filter(
                                                        partitionInfo -> {
                                                            if (kafkaSourceConfig != null
                                                                    && kafkaSourceConfig
                                                                            .isIgnoreNoLeaderPartition()
                                                                    && partitionInfo.leader()
                                                                            == null) {
                                                                log.warn(
                                                                        "Partition {} of topic {} has no leader, skipping due to ignore_no_leader_partition=true.",
                                                                        partitionInfo.partition(),
                                                                        t.name());
                                                                return false;
                                                            }
                                                            return true;
                                                        })
                                                .map(
                                                        p ->
                                                                new TopicPartition(
                                                                        t.name(), p.partition())))
                        .collect(Collectors.toSet());
        Map<TopicPartition, Long> latestOffsets = listOffsets(partitions, OffsetSpec.latest());
        return partitions.stream()
                .map(
                        partition -> {
                            // Obtain the corresponding topic TablePath from kafka topic
                            TablePath tablePath = topicMappingTablePathMap.get(partition.topic());

View on GitHub (pinned to cf67b549a7)