apache/iceberg · warning

Rebalance may have occurred, partition

Error message

Rebalance may have occurred, partition {} lost before seeking

What it means

While seeking each partition to its last committed offset, the consumer threw IllegalStateException, which Kafka raises when the consumer no longer belongs to the group for that partition — typically because a rebalance revoked the partition before the seek. The exception is caught and logged per-partition so the remaining partitions can still be positioned.

Solutions

  1. No action needed for correctness — after the rebalance the framework repositions partitions from committed offsets anyway; this warning is informational.
  2. If frequent, increase max.poll.interval.ms and reduce per-poll work so rebalances do not interrupt seekToLastCommittedOffsets.
  3. Check for repeated consumer group membership churn (flapping workers, session.timeout.ms too low) and stabilize worker count.
  4. Verify records are not duplicated after rebalance in downstream tables; if duplicates appear, review the offset-commit strategy (kafka.commit.group-id / offset commit settings).
Defensive patterns

Strategy: try-catch

Try / catch

try {
  consumer.seek(topicPartition, offsetAndMetadata.offset());
} catch (IllegalStateException e) {
  // partition revoked by rebalance; skip — framework repositions after reassignment
  LOG.warn("Rebalance may have occurred, partition {} lost before seeking", topicPartition, e);
}

Prevention

When it happens

Trigger: KafkaUtils.seekToLastCommittedOffsets calls consumer.seek(topicPartition, offset) and the partition was revoked in a concurrent rebalance (consumer is no longer assigned that partition).

Common situations: Consumer group rebalance triggered while the channel was starting up (member join/leave, session timeout, max.poll.interval exceeded); long initialization between poll and seek allowing the rebalance to revoke partitions.

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/7ac3e4bd24d9adf3. Report an issue: GitHub.

Appendix: source

Thrown at kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java:78

  static void seekToLastCommittedOffsets(SinkTaskContext context) {
    Consumer<byte[], byte[]> consumer = kafkaConsumer(context);
    if (consumer == null) {
      return;
    }

    Map<TopicPartition, OffsetAndMetadata> committedOffsets =
        consumer.committed(consumer.assignment());
    if (committedOffsets == null || committedOffsets.isEmpty()) {
      return;
    }

    committedOffsets.forEach(
        (topicPartition, offsetAndMetadata) -> {
          if (offsetAndMetadata != null) {
            try {
              consumer.seek(topicPartition, offsetAndMetadata.offset());
            } catch (IllegalStateException e) {
              LOG.warn(
                  "Rebalance may have occurred, partition {} lost before seeking",
                  topicPartition,
                  e);
            }
          }
        });
  }

  @SuppressWarnings("unchecked")
  private static Consumer<byte[], byte[]> kafkaConsumer(SinkTaskContext context) {
    String contextClassName = context.getClass().getName();
    try {
      return ((Consumer<byte[], byte[]>)
          DynFields.builder().hiddenImpl(CONTEXT_CLASS_NAME, "consumer").build(context).get());
    } catch (Exception e) {
      throw new ConnectException(
          "Unable to retrieve consumer from context: " + contextClassName, e);
    }

View on GitHub (pinned to 86d9c8fc54)