apache/kafka · error · IllegalArgumentException

Topic must be non-null.

Error message

Topic must be non-null.

What it means

Thrown by ConsumerRecords.records(String topic) when the topic argument is null. The method filters the records map by topic equality, and a null topic would never match (and would be ambiguous), so Kafka rejects it up front. Pass the topic you want records for; if you want everything, use the iterator or .count() instead.

Solutions

  1. Validate the topic string is non-null before calling records(topic); guard with Objects.requireNonNull(topic, "topic").
  2. Load the topic name from a config key with a sane default and fail fast at startup if missing.
  3. If you want all records, iterate ConsumerRecords directly instead of calling records(topic).

Example fix

// before
Iterable<ConsumerRecord<K,V>> recs = polled.records(topicFromConfig); // topicFromConfig == null

// after
String topic = Objects.requireNonNull(topicFromConfig, "topic must be configured");
Iterable<ConsumerRecord<K,V>> recs = polled.records(topic);
Defensive patterns

Strategy: validation

Validate before calling

String safeTopic = Objects.requireNonNull(topic, "topic must be non-null before records(topic)");
return polled.records(safeTopic);

Type guard

static boolean isQueryableTopic(String t) {
    return t != null && !t.isEmpty();
}

Try / catch

try {
    return polled.records(topic);
} catch (IllegalArgumentException e) {
    if ("Topic must be non-null.".equals(e.getMessage())) {
        return Collections.emptyList(); // or surface as a configuration error
    }
    throw e;
}

Prevention

When it happens

Trigger: Calling consumer.poll(...).records(topic) where topic is a variable that resolved to null; passing a topic name pulled from a config/property that was unset; looping over a collection of topic names where one entry is null.

Common situations: Apps that read target topics from external config and forget a default; helper methods that accept a nullable topic parameter; tests with placeholder topic variables.

Related errors


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

Appendix: source

Thrown at clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerRecords.java:122

            // every call. A time based approach is used to avoid this. See KAFKA-20660 for more details.
            if (now - lastLog >= TAINT_LOG_INTERVAL_NS && TAINTED_NEXT_OFFSETS_LAST_LOG_NS.compareAndSet(lastLog, now)) {
                log.error("ConsumerRecords#nextOffsets() returned empty because this instance was built with the " +
                        "deprecated ConsumerRecords(Map) constructor (see KIP-1094), which does not supply next offsets. " +
                        "Downstream logic that relies on these offsets to advance the consumer's committed position " +
                        "(for example, Kafka Streams under exactly-once semantics) will be unable to commit, leading to " +
                        "reprocessing. Update the interceptor or wrapper that constructed it to use the " +
                        "ConsumerRecords(Map, Map) constructor that supplies next offsets.");
            }
        }
        return nextOffsets;
    }

    /**
     * Get just the records for the given topic
     */
    public Iterable<ConsumerRecord<K, V>> records(String topic) {
        if (topic == null)
            throw new IllegalArgumentException("Topic must be non-null.");
        List<List<ConsumerRecord<K, V>>> recs = new ArrayList<>();
        for (Map.Entry<TopicPartition, List<ConsumerRecord<K, V>>> entry : records.entrySet()) {
            if (entry.getKey().topic().equals(topic))
                recs.add(entry.getValue());
        }
        return new ConcatenatedIterable<>(recs);
    }

    /**
     * Get the partitions which have records contained in this record set.
     * @return The set of partitions with data in this record set (may be empty if no data was returned)
     */
    public Set<TopicPartition> partitions() {
        return Collections.unmodifiableSet(records.keySet());
    }

    @Override
    public Iterator<ConsumerRecord<K, V>> iterator() {

View on GitHub (pinned to 996fb4585a)