{"record":{"id":"5170e2fceb970530","repo":"apache/kafka","slug":"consumer-is-not-subscribed-to-any-topics-or-assign","errorCode":null,"errorMessage":"Consumer is not subscribed to any topics or assigned any partitions","messagePattern":"Consumer is not subscribed to any topics or assigned any partitions","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java","lineNumber":939,"sourceCode":"     *             partitions to consume from or an unexpected error occurred\n     * @throws org.apache.kafka.clients.consumer.OffsetOutOfRangeException if the fetch position of the consumer is\n     *             out of range and no offset reset policy is configured.\n     * @throws org.apache.kafka.common.errors.TopicAuthorizationException if the consumer is not authorized to read\n     *             from a partition\n     * @throws org.apache.kafka.common.errors.SerializationException if the fetched records cannot be deserialized\n     * @throws org.apache.kafka.common.errors.UnsupportedAssignorException if the `group.remote.assignor` configuration\n     *             is set to an assignor that is not available on the broker.\n     */\n    @Override\n    public ConsumerRecords<K, V> poll(final Duration timeout) {\n        Timer timer = time.timer(timeout);\n\n        acquireAndEnsureOpen();\n        try {\n            kafkaConsumerMetrics.recordPollStart(timer.currentTimeMs());\n\n            if (subscriptions.hasNoSubscriptionOrUserAssignment()) {\n                throw new IllegalStateException(\"Consumer is not subscribed to any topics or assigned any partitions\");\n            }\n\n            // This distinguishes the first pass of the inner do/while loop from subsequent passes for the\n            // inflight poll event logic.\n            boolean firstPass = true;\n\n            do {\n                // We must not allow wake-ups between polling for fetches and returning the records.\n                // If the polled fetches are not empty the consumed position has already been updated in the polling\n                // of the fetches. A wakeup between returned fetches and returning records would lead to never\n                // returning the records in the fetches. Thus, we trigger a possible wake-up before we poll fetches.\n                wakeupTrigger.maybeTriggerWakeup();\n\n                checkInflightPoll(timer, firstPass);\n                firstPass = false;\n                final Fetch<K, V> fetch = pollForFetches(timer);\n                if (!fetch.isEmpty()) {\n                    // before returning the fetched records, we can send off the next round of fetches","sourceCodeStart":921,"sourceCodeEnd":957,"githubUrl":"https://github.com/apache/kafka/blob/996fb4585aa1bcc8980b0e1b8d6b168b986cd979/clients/src/main/java/org/apache/kafka/clients/consumer/internals/AsyncKafkaConsumer.java#L921-L957","documentation":"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.","triggerScenarios":"Invoking consumer.poll(...) before any consumer.subscribe(...) or consumer.assign(...), or after unsubscribe() with no re-subscription.","commonSituations":"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).","solutions":["Call subscribe(Collections.singleton(topic)) or assign(Collections.singleton(tp)) before the first poll.","If running conditionally, guard poll with a check that subscription/assignment is non-empty.","After unsubscribe(), re-subscribe or re-assign before the next poll."],"exampleFix":"// before\nKafkaConsumer<String,String> c = new KafkaConsumer<>(props);\nc.poll(Duration.ofMillis(100)); // -> IllegalStateException\n\n// after\nKafkaConsumer<String,String> c = new KafkaConsumer<>(props);\nc.subscribe(Collections.singleton(\"t\"));\nc.poll(Duration.ofMillis(100));","handlingStrategy":"validation","validationCode":"if (consumer.subscription().isEmpty() && consumer.assignment().isEmpty()) {\n    throw new IllegalStateException(\"subscribe or assign before poll\");\n}\nconsumer.poll(timeout);","typeGuard":"null","tryCatchPattern":"null  // IllegalStateException here is a programming error; fix the call ordering","preventionTips":["Always subscribe/assign in the same setup method that creates the consumer.","After unsubscribe(), re-subscribe before the next poll.","In tests, do this in @BeforeEach."],"tags":["consumer","poll","subscription","lifecycle","kafka-clients"],"backgroundTag":null,"analyzedSha":"996fb4585aa1bcc8980b0e1b8d6b168b986cd979","analyzedAt":"2026-08-11T22:03:28.655Z","contentChangedAt":null,"schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}