apache/kafka · error · IllegalArgumentException

Topic partition was not included in the request

Error message

Topic partition {partition} was not included in the request

What it means

Thrown by DescribeProducersResult.partitionResult() when the given TopicPartition has no corresponding entry in the futures map captured at construction. This means the partition was not part of the original Admin.describeProducersShuffleStrategies / describeProducers request.

Solutions

  1. Only call partitionResult with TopicPartitions that were passed to the original describeProducers request.
  2. Cache the original partition collection and check membership before lookup.
  3. Use the all() future if you need results for every requested partition without per-partition lookup.

Example fix

// before
result.partitionResult(new TopicPartition(topic, p));

// after
if (requestedPartitions.contains(tp)) {
    result.partitionResult(tp);
}
Defensive patterns

Strategy: validation

Validate before calling

Set<TopicPartition> requested = ...;
if (!requested.contains(tp)) {
    throw new IllegalArgumentException(tp + " not in original request");
}
result.partitionResult(tp);

Type guard

boolean wasRequested(Set<TopicPartition> requested, TopicPartition tp) {
    return requested.contains(tp);
}

Prevention

When it happens

Trigger: Calling partitionResult(tp) where futures.get(tp) returns null — tp was not requested in the describeProducers call.

Common situations: Caller queries a partition that was not in the original request, or constructs a TopicPartition with a mismatched topic name or partition number.

Related errors


AI-assisted analysis of apache/kafka@996fb4585a (2026-08-11). Data as JSON: /api/errors/5afa3b32ff2340ed. Report an issue: GitHub.

Appendix: source

Thrown at clients/src/main/java/org/apache/kafka/clients/admin/DescribeProducersResult.java:42

import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutionException;

@InterfaceAudience.Public
public class DescribeProducersResult {

    private final Map<TopicPartition, KafkaFuture<PartitionProducerState>> futures;

    DescribeProducersResult(Map<TopicPartition, KafkaFuture<PartitionProducerState>> futures) {
        this.futures = futures;
    }

    public KafkaFuture<PartitionProducerState> partitionResult(final TopicPartition partition) {
        KafkaFuture<PartitionProducerState> future = futures.get(partition);
        if (future == null) {
            throw new IllegalArgumentException("Topic partition " + partition +
                " was not included in the request");
        }
        return future;
    }

    public KafkaFuture<Map<TopicPartition, PartitionProducerState>> all() {
        return KafkaFuture.allOf(futures.values().toArray(new KafkaFuture<?>[0]))
            .thenApply(nil -> {
                Map<TopicPartition, PartitionProducerState> results = new HashMap<>(futures.size());
                for (Map.Entry<TopicPartition, KafkaFuture<PartitionProducerState>> entry : futures.entrySet()) {
                    try {
                        results.put(entry.getKey(), entry.getValue().get());
                    } catch (InterruptedException | ExecutionException e) {
                        // This should be unreachable, because allOf ensured that all the futures completed successfully.
                        throw new KafkaException(e);
                    }
                }
                return results;

View on GitHub (pinned to 996fb4585a)