{"record":{"id":"1010850ea1897a94","repo":"apache/iceberg","slug":"already-closed-files-for-partition-path","errorCode":null,"errorMessage":"Already closed files for partition: <path>","messagePattern":"Already closed files for partition: <path>","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"core/src/main/java/org/apache/iceberg/io/PartitionedWriter.java","lineNumber":84,"sourceCode":"   */\n  protected abstract PartitionKey partition(T row);\n\n  @Override\n  public void write(T row) throws IOException {\n    PartitionKey key = partition(row);\n\n    if (!key.equals(currentKey)) {\n      if (currentKey != null) {\n        // if the key is null, there was no previous current key and current writer.\n        currentWriter.close();\n        completedPartitions.add(currentKey);\n      }\n\n      if (completedPartitions.contains(key)) {\n        // if rows are not correctly grouped, detect and fail the write\n        PartitionKey existingKey = Iterables.find(completedPartitions, key::equals, null);\n        LOG.warn(\"Duplicate key: {} == {}\", existingKey, key);\n        throw new IllegalStateException(\"Already closed files for partition: \" + key.toPath());\n      }\n\n      currentKey = key.copy();\n      currentWriter = new RollingFileWriter(currentKey);\n    }\n\n    currentWriter.write(row);\n  }\n\n  @Override\n  public void close() throws IOException {\n    if (currentWriter != null) {\n      currentWriter.close();\n    }\n  }\n}\n","sourceCodeStart":66,"sourceCodeEnd":101,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/core/src/main/java/org/apache/iceberg/io/PartitionedWriter.java#L66-L101","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":"// before\ndf.writeTo(\"db.table\").using(\"iceberg\").append(); // rows not clustered by partition\n\n// after\ndf.repartitionByRange(col(\"dt\"), col(\"country\")) // match partition spec\n  .sortWithinPartitions(\"dt\", \"country\")\n  .writeTo(\"db.table\").using(\"iceberg\").append();\n// or set table property: 'write.fanout.enabled'='true'","handlingStrategy":"validation","validationCode":"// Scala/Spark: ensure rows are clustered by partition columns before writing\ndf.select(partitionCols: _*).distinct() // sanity check exists\nassert(df.queryExecution.executedPlan.outputPartitioning.satisfies(\n  Distribution.createOrderedDistribution(partitionExprs)),\n  \"Input must be partitioned/sorted by partition columns\")","typeGuard":null,"tryCatchPattern":"try {\n  writer.write(row);\n} catch (IllegalStateException e) {\n  if (e.getMessage().startsWith(\"Already closed files\")) {\n    // abort task and re-execute with fanout.enabled=true or sorted input\n  }\n  throw e;\n}","preventionTips":["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)"],"tags":["unsorted-input","partitioned-write","data-corruption-prevention"],"backgroundTag":"invalid-state-transition","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"}