{"record":{"id":"b8c23ba5a7de0020","repo":"apache/iceberg","slug":"unsupported-table-change-addwatermark-b8c23b","errorCode":null,"errorMessage":"Unsupported table change: AddWatermark.","messagePattern":"Unsupported table change: AddWatermark\\.","errorType":"exception","errorClass":"UnsupportedOperationException","httpStatus":null,"severity":"error","filePath":"flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/util/FlinkAlterTableUtil.java","lineNumber":135,"sourceCode":"  /**\n   * Applies a list of Flink table changes to an {@link UpdateSchema} operation.\n   *\n   * @param pendingUpdate an uncommitted UpdateSchema operation to configure\n   * @param schemaChanges a list of Flink table changes\n   */\n  public static void applySchemaChanges(\n      UpdateSchema pendingUpdate, List<TableChange> schemaChanges) {\n    for (TableChange change : schemaChanges) {\n      if (change instanceof TableChange.AddColumn) {\n        applyAddColumn(pendingUpdate, (TableChange.AddColumn) change);\n      } else if (change instanceof TableChange.ModifyColumn) {\n        TableChange.ModifyColumn modifyColumn = (TableChange.ModifyColumn) change;\n        applyModifyColumn(pendingUpdate, modifyColumn);\n      } else if (change instanceof TableChange.DropColumn) {\n        TableChange.DropColumn dropColumn = (TableChange.DropColumn) change;\n        pendingUpdate.deleteColumn(dropColumn.getColumnName());\n      } else if (change instanceof TableChange.AddWatermark) {\n        throw new UnsupportedOperationException(\"Unsupported table change: AddWatermark.\");\n      } else if (change instanceof TableChange.ModifyWatermark) {\n        throw new UnsupportedOperationException(\"Unsupported table change: ModifyWatermark.\");\n      } else if (change instanceof TableChange.DropWatermark) {\n        throw new UnsupportedOperationException(\"Unsupported table change: DropWatermark.\");\n      } else if (change instanceof TableChange.AddUniqueConstraint) {\n        TableChange.AddUniqueConstraint addPk = (TableChange.AddUniqueConstraint) change;\n        applyUniqueConstraint(pendingUpdate, addPk.getConstraint());\n      } else if (change instanceof TableChange.ModifyUniqueConstraint) {\n        TableChange.ModifyUniqueConstraint modifyPk = (TableChange.ModifyUniqueConstraint) change;\n        applyUniqueConstraint(pendingUpdate, modifyPk.getNewConstraint());\n      } else if (change instanceof TableChange.DropConstraint) {\n        throw new UnsupportedOperationException(\"Unsupported table change: DropConstraint.\");\n      } else {\n        throw new UnsupportedOperationException(\"Cannot apply unknown table change: \" + change);\n      }\n    }\n  }\n","sourceCodeStart":117,"sourceCodeEnd":153,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/util/FlinkAlterTableUtil.java#L117-L153","documentation":"FlinkAlterTableUtil.applySchemaChanges translates Flink TableChange objects into Iceberg schema updates. Watermark changes are not supported by Iceberg tables, so TableChange.AddWatermark throws UnsupportedOperationException.","triggerScenarios":"Executing Flink SQL 'ALTER TABLE ... ADD WATERMARK ...' (or calling the API with a TableChange.AddWatermark) against an Iceberg catalog table.","commonSituations":"Porting Flink SQL DDL written for other connectors to Iceberg; attempts to add event-time watermarks to an Iceberg table schema.","solutions":["Remove the ADD WATERMARK clause from the DDL — watermarks belong in the query/plan, not the Iceberg table","Define watermark logic in the Flink job's source/windowing instead of the table","If watermark semantics must persist, consider Iceberg watermark-related table properties at the engine level (not via schema change)"],"exampleFix":"// before\n// ALTER TABLE iceberg_table ADD WATERMARK FOR rowtime AS rowtime - INTERVAL '5' SECOND;\n// after: apply watermark at query time, or drop the clause entirely\nINSERT INTO iceberg_table SELECT ... ; // watermark handled in the pipeline","handlingStrategy":"validation","validationCode":"boolean hasWatermarkChange = java.util.Arrays.stream(changes)\n    .anyMatch(c -> c instanceof TableChange.AddWatermark\n        || c instanceof TableChange.ModifyWatermark\n        || c instanceof TableChange.DropWatermark);\nif (hasWatermarkChange) {\n  throw new IllegalArgumentException(\"Iceberg tables do not support watermark table changes\");\n}","typeGuard":null,"tryCatchPattern":"try {\n  FlinkAlterTableUtil.applySchemaChanges(update, changes);\n} catch (UnsupportedOperationException e) {\n  LOG.warn(\"Ignoring unsupported change for Iceberg: {}\", e.getMessage());\n  // or rethrow if the change is essential\n}","preventionTips":["Never use ADD WATERMARK DDL on Iceberg tables","Handle watermarks in the Flink pipeline (source/window operators), not table schema","Filter watermark changes out of generic DDL-sync tooling before it touches Iceberg"],"tags":["flink","ddl","unsupported-operation","watermark"],"backgroundTag":"unsupported-operation","analyzedSha":"86d9c8fc543e7c56c9f624eb725f76c9baff9570","analyzedAt":"2026-09-12T00:46:39.097Z","contentChangedAt":"2026-09-12T00:46:39.097Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}