apache/kafka · error · IllegalArgumentException

You can only check the position for partitions assigned to…

Error message

You can only check the position for partitions assigned to this consumer.

What it means

Thrown by MockConsumer.position(TopicPartition) when the partition is not in the consumer's current assignment. Position only has meaning for partitions the consumer is actively consuming; asking for a position on an unassigned partition is a programming error in the test. The guard mirrors the real consumer's behavior.

Solutions

  1. Ensure the TopicPartition is in mock.assignment() before calling position().
  2. Call assign() or simulate rebalance so the partition is assigned.
  3. Derive tp from the same constant used during assign to avoid drift.

Example fix

// before
TopicPartition tp = new TopicPartition("orders", 0);
long pos = mockConsumer.position(tp); // not assigned

// after
mockConsumer.assign(Set.of(tp));
mockConsumer.updateBeginningOffsets(Map.of(tp, 0L));
long pos = mockConsumer.position(tp);
Defensive patterns

Strategy: validation

Validate before calling

if (!mockConsumer.assignment().contains(tp))
    throw new IllegalStateException("position() requires " + tp + " to be assigned");
return mockConsumer.position(tp);

Type guard

static boolean canQueryPosition(MockConsumer<?,?> m, TopicPartition tp) {
    return m.assignment().contains(tp);
}

Try / catch

try {
    return mockConsumer.position(tp);
} catch (IllegalArgumentException e) {
    if ("You can only check the position for partitions assigned to this consumer.".equals(e.getMessage())) {
        mockConsumer.assign(union(mockConsumer.assignment(), Set.of(tp)));
        return mockConsumer.position(tp);
    }
    throw e;
}

Prevention

When it happens

Trigger: Calling position(tp) before assign(); tp with wrong topic name or partition number; calling position after revoke but before re-assign.

Common situations: Test setup ordering bug (position before assign); stale TopicPartition references after a simulated rebalance; refactor that changed partition counts.

Related errors


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

Appendix: source

Thrown at clients/src/main/java/org/apache/kafka/clients/consumer/MockConsumer.java:489

    public synchronized Map<TopicPartition, OffsetAndMetadata> committed(final Set<TopicPartition> partitions) {
        ensureNotClosed();

        return partitions.stream()
            .filter(committed::containsKey)
            .collect(Collectors.toMap(tp -> tp, tp -> subscriptions.isAssigned(tp) ?
                committed.get(tp) : new OffsetAndMetadata(0)));
    }

    @Override
    public synchronized Map<TopicPartition, OffsetAndMetadata> committed(final Set<TopicPartition> partitions, final Duration timeout) {
        return committed(partitions);
    }

    @Override
    public synchronized long position(TopicPartition partition) {
        ensureNotClosed();
        if (!this.subscriptions.isAssigned(partition))
            throw new IllegalArgumentException("You can only check the position for partitions assigned to this consumer.");
        SubscriptionState.FetchPosition position = this.subscriptions.position(partition);
        if (position == null) {
            updateFetchPosition(partition);
            position = this.subscriptions.position(partition);
        }
        return position.offset;
    }

    @Override
    public synchronized long position(TopicPartition partition, final Duration timeout) {
        return position(partition);
    }

    @Override
    public synchronized void seekToBeginning(Collection<TopicPartition> partitions) {
        ensureNotClosed();
        subscriptions.requestOffsetReset(partitions, AutoOffsetResetStrategy.EARLIEST);
    }

View on GitHub (pinned to 996fb4585a)