{"record":{"id":"7ac3e4bd24d9adf3","repo":"apache/iceberg","slug":"rebalance-may-have-occurred-partition-lost-bef","errorCode":null,"errorMessage":"Rebalance may have occurred, partition {} lost before seeking","messagePattern":"Rebalance may have occurred, partition (.+?) lost before seeking","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java","lineNumber":78,"sourceCode":"  static void seekToLastCommittedOffsets(SinkTaskContext context) {\n    Consumer<byte[], byte[]> consumer = kafkaConsumer(context);\n    if (consumer == null) {\n      return;\n    }\n\n    Map<TopicPartition, OffsetAndMetadata> committedOffsets =\n        consumer.committed(consumer.assignment());\n    if (committedOffsets == null || committedOffsets.isEmpty()) {\n      return;\n    }\n\n    committedOffsets.forEach(\n        (topicPartition, offsetAndMetadata) -> {\n          if (offsetAndMetadata != null) {\n            try {\n              consumer.seek(topicPartition, offsetAndMetadata.offset());\n            } catch (IllegalStateException e) {\n              LOG.warn(\n                  \"Rebalance may have occurred, partition {} lost before seeking\",\n                  topicPartition,\n                  e);\n            }\n          }\n        });\n  }\n\n  @SuppressWarnings(\"unchecked\")\n  private static Consumer<byte[], byte[]> kafkaConsumer(SinkTaskContext context) {\n    String contextClassName = context.getClass().getName();\n    try {\n      return ((Consumer<byte[], byte[]>)\n          DynFields.builder().hiddenImpl(CONTEXT_CLASS_NAME, \"consumer\").build(context).get());\n    } catch (Exception e) {\n      throw new ConnectException(\n          \"Unable to retrieve consumer from context: \" + contextClassName, e);\n    }","sourceCodeStart":60,"sourceCodeEnd":96,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/KafkaUtils.java#L60-L96","documentation":"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.","triggerScenarios":"KafkaUtils.seekToLastCommittedOffsets calls consumer.seek(topicPartition, offset) and the partition was revoked in a concurrent rebalance (consumer is no longer assigned that partition).","commonSituations":"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.","solutions":["No action needed for correctness — after the rebalance the framework repositions partitions from committed offsets anyway; this warning is informational.","If frequent, increase max.poll.interval.ms and reduce per-poll work so rebalances do not interrupt seekToLastCommittedOffsets.","Check for repeated consumer group membership churn (flapping workers, session.timeout.ms too low) and stabilize worker count.","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)."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"try {\n  consumer.seek(topicPartition, offsetAndMetadata.offset());\n} catch (IllegalStateException e) {\n  // partition revoked by rebalance; skip — framework repositions after reassignment\n  LOG.warn(\"Rebalance may have occurred, partition {} lost before seeking\", topicPartition, e);\n}","preventionTips":["Increase max.poll.interval.ms if seek positioning regularly overlaps rebalances.","Stabilize consumer group membership (avoid flapping workers, tune session.timeout.ms).","Treat the warning as informational — correctness is preserved by committed offsets.","Monitor for downstream duplicates after rebalances to validate offset-commit settings."],"tags":["kafka","consumer","rebalance","seek"],"backgroundTag":"consumer-rebalance-revoked-partition","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}