apache/kafka · error · IllegalStateException

Consumer is not subscribed to any topics or assigned any…

Error message

Consumer is not subscribed to any topics or assigned any partitions

What it means

AsyncKafkaConsumer.poll guards against a no-op poll: subscriptions.hasNoSubscriptionOrUserAssignment() returns true when neither subscribe() nor assign() has been called. Calling poll in that state throws IllegalStateException because there is nothing to fetch and the call would misleadingly return empty records forever. This is checked after acquireAndEnsureOpen(), so the consumer must also be open.

Solutions

  1. Call subscribe(Collections.singleton(topic)) or assign(Collections.singleton(tp)) before the first poll.
  2. If running conditionally, guard poll with a check that subscription/assignment is non-empty.
  3. After unsubscribe(), re-subscribe or re-assign before the next poll.

Example fix

// before
KafkaConsumer<String,String> c = new KafkaConsumer<>(props);
c.poll(Duration.ofMillis(100)); // -> IllegalStateException

// after
KafkaConsumer<String,String> c = new KafkaConsumer<>(props);
c.subscribe(Collections.singleton("t"));
c.poll(Duration.ofMillis(100));
Defensive patterns

Strategy: validation

Validate before calling

if (consumer.subscription().isEmpty() && consumer.assignment().isEmpty()) {
    throw new IllegalStateException("subscribe or assign before poll");
}
consumer.poll(timeout);

Type guard

null

Try / catch

null  // IllegalStateException here is a programming error; fix the call ordering

Prevention

When it happens

Trigger: Invoking consumer.poll(...) before any consumer.subscribe(...) or consumer.assign(...), or after unsubscribe() with no re-subscription.

Common situations: Test or application wiring where subscribe() is conditional and skipped; calling poll in a @BeforeEach before the test body subscribes; misordered lifecycle (poll fired from a scheduler before setup completes).

Related errors


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

Appendix: source

Thrown at clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java:939

     *             partitions to consume from or an unexpected error occurred
     * @throws org.apache.kafka.clients.consumer.OffsetOutOfRangeException if the fetch position of the consumer is
     *             out of range and no offset reset policy is configured.
     * @throws org.apache.kafka.common.errors.TopicAuthorizationException if the consumer is not authorized to read
     *             from a partition
     * @throws org.apache.kafka.common.errors.SerializationException if the fetched records cannot be deserialized
     * @throws org.apache.kafka.common.errors.UnsupportedAssignorException if the `group.remote.assignor` configuration
     *             is set to an assignor that is not available on the broker.
     */
    @Override
    public ConsumerRecords<K, V> poll(final Duration timeout) {
        Timer timer = time.timer(timeout);

        acquireAndEnsureOpen();
        try {
            kafkaConsumerMetrics.recordPollStart(timer.currentTimeMs());

            if (subscriptions.hasNoSubscriptionOrUserAssignment()) {
                throw new IllegalStateException("Consumer is not subscribed to any topics or assigned any partitions");
            }

            // This distinguishes the first pass of the inner do/while loop from subsequent passes for the
            // inflight poll event logic.
            boolean firstPass = true;

            do {
                // We must not allow wake-ups between polling for fetches and returning the records.
                // If the polled fetches are not empty the consumed position has already been updated in the polling
                // of the fetches. A wakeup between returned fetches and returning records would lead to never
                // returning the records in the fetches. Thus, we trigger a possible wake-up before we poll fetches.
                wakeupTrigger.maybeTriggerWakeup();

                checkInflightPoll(timer, firstPass);
                firstPass = false;
                final Fetch<K, V> fetch = pollForFetches(timer);
                if (!fetch.isEmpty()) {
                    // before returning the fetched records, we can send off the next round of fetches

View on GitHub (pinned to 996fb4585a)