apache/iceberg · error · IllegalStateException

Already closed files for partition

Error message

Already closed files for partition: <path>

What it means

PartitionedWriter detects that rows for a partition key arrive after the writer for that partition has already been closed. Iceberg requires all rows of a partition to be grouped contiguously; when a partition reappears after another partition started, the already-closed data files would be duplicated or corrupted, so the write fails fast with IllegalStateException.

Solutions

  1. Sort or repartition the input rows by all partition columns before writing.
  2. Set write.fanout.enabled=true on the table properties if input ordering cannot be guaranteed.
  3. Verify the partition spec columns match the columns actually used for upstream clustering.
  4. Check engine-level shuffle/repartition configuration (e.g. Spark distribute-by) covers every partition column.

Example fix

// before
df.writeTo("db.table").using("iceberg").append(); // rows not clustered by partition

// after
df.repartitionByRange(col("dt"), col("country")) // match partition spec
  .sortWithinPartitions("dt", "country")
  .writeTo("db.table").using("iceberg").append();
// or set table property: 'write.fanout.enabled'='true'
Defensive patterns

Strategy: validation

Validate before calling

// Scala/Spark: ensure rows are clustered by partition columns before writing
df.select(partitionCols: _*).distinct() // sanity check exists
assert(df.queryExecution.executedPlan.outputPartitioning.satisfies(
  Distribution.createOrderedDistribution(partitionExprs)),
  "Input must be partitioned/sorted by partition columns")

Try / catch

try {
  writer.write(row);
} catch (IllegalStateException e) {
  if (e.getMessage().startsWith("Already closed files")) {
    // abort task and re-execute with fanout.enabled=true or sorted input
  }
  throw e;
}

Prevention

When it happens

Trigger: Calling PartitionedWriter.write(row) with a row whose PartitionKey equals a key in completedPartitions, i.e. rows for the same partition not sorted/clustered contiguously in the input stream.

Common situations: Upstream sort/cluster on partition columns dropped or changed; a shuffle repartitioned by the wrong key; multiple map tasks interleaving partitions after a merge; fanout disabled while input ordering is not guaranteed by partition columns.

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.


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

Appendix: source

Thrown at core/src/main/java/org/apache/iceberg/io/PartitionedWriter.java:84

   */
  protected abstract PartitionKey partition(T row);

  @Override
  public void write(T row) throws IOException {
    PartitionKey key = partition(row);

    if (!key.equals(currentKey)) {
      if (currentKey != null) {
        // if the key is null, there was no previous current key and current writer.
        currentWriter.close();
        completedPartitions.add(currentKey);
      }

      if (completedPartitions.contains(key)) {
        // if rows are not correctly grouped, detect and fail the write
        PartitionKey existingKey = Iterables.find(completedPartitions, key::equals, null);
        LOG.warn("Duplicate key: {} == {}", existingKey, key);
        throw new IllegalStateException("Already closed files for partition: " + key.toPath());
      }

      currentKey = key.copy();
      currentWriter = new RollingFileWriter(currentKey);
    }

    currentWriter.write(row);
  }

  @Override
  public void close() throws IOException {
    if (currentWriter != null) {
      currentWriter.close();
    }
  }
}

View on GitHub (pinned to 86d9c8fc54)