apache/iceberg · warning

Close unexpectedly called on committer

Error message

Close unexpectedly called on committer {} without partition assignment

What it means

CommitterImpl.close was called before the task was ever initialized with a partition assignment. This is defensively tolerated (worker is still stopped) and logged as a warning, because closing an uninitialized committer indicates an unexpected SourceTask lifecycle from the Connect framework.

Solutions

  1. Usually safe to ignore — close() stops the worker and returns; no state corruption occurs.
  2. If it appears at every startup, check the task log above for a start()-time exception (e.g. catalog/auth failure) that prevented initialization.
  3. Validate the connector configuration (catalog properties, topics, table settings) so the task initializes on first attempt.
  4. Upgrade the connector if your Connect runtime invokes close() during failed startup in a way that repeatedly triggers this.
Defensive patterns

Strategy: validation

Validate before calling

// validate connector config up front so the task initializes on first attempt
// (catalog properties, topics, routes, table identifiers)

Prevention

When it happens

Trigger: Kafka Connect calls close() on the committer task that was never opened()/initialized — e.g. task failed during start before assignment, or immediate task shutdown after configuration.

Common situations: Task startup aborted early due to config errors, so close follows without initialization; rapid task reconfiguration by the Connect worker; framework edge cases during worker shutdown.

Understand the failure class

Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.

Related errors


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

Appendix: source

Thrown at kafka-connect/kafka-connect/src/main/java/org/apache/iceberg/connect/channel/CommitterImpl.java:167

      startCoordinator();
    }
  }

  @Override
  public void stop() {
    throw new UnsupportedOperationException(
        "The method stop() is deprecated and will be removed in 2.0.0. "
            + "Use stop(Collection<TopicPartition>) instead.");
  }

  @Override
  public void close(Collection<TopicPartition> closedPartitions) {
    // Always try to stop the worker to avoid duplicates.
    stopWorker();

    // Defensive: close called without prior initialization (should not happen).
    if (!isInitialized.get()) {
      LOG.warn("Close unexpectedly called on committer {} without partition assignment", taskId);
      return;
    }

    // Empty partitions → task was stopped explicitly. Stop coordinator if running.
    if (closedPartitions.isEmpty()) {
      LOG.info("Committer {} stopped. Closing coordinator.", taskId);
      stopCoordinator();
      return;
    }

    // Normal close: if leader partition is lost, stop coordinator.
    if (hasLeaderPartition(closedPartitions)) {
      LOG.info("Committer {} lost leader partition. Stopping coordinator.", taskId);
      stopCoordinator();
    }

    // Reset offsets to last committed to avoid data loss.
    LOG.info("Seeking to last committed offsets for worker {}.", taskId);

View on GitHub (pinned to 86d9c8fc54)