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
- Sort or repartition the input rows by all partition columns before writing.
- Set write.fanout.enabled=true on the table properties if input ordering cannot be guaranteed.
- Verify the partition spec columns match the columns actually used for upstream clustering.
- 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
- Enable write.fanout.enabled=true for tables fed by unsorted streams
- Always sortWithinPartitions/repartition by the exact partition spec columns
- Keep upstream clustering logic and partition spec changes in sync (review spec evolution)
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)