{"record":{"id":"8c6d4fb7a754ea2e","repo":"apache/iceberg","slug":"exception-processing-split-at","errorCode":null,"errorMessage":"Exception processing split {} at {}","messagePattern":"Exception processing split (.+?) at (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/TableReader.java","lineNumber":101,"sourceCode":"            .counter(TableMaintenanceMetrics.ERROR_COUNTER);\n    this.rowDataReaderFunction =\n        new MetaDataReaderFunction(\n            new Configuration(),\n            metaTable.schema(),\n            projectedSchema,\n            metaTable.io(),\n            metaTable.encryption());\n    this.splitSerializer = new IcebergSourceSplitSerializer(scanContext.caseSensitive());\n  }\n\n  @Override\n  public void processElement(\n      MetadataTablePlanner.SplitInfo splitInfo, Context ctx, Collector<R> out) throws Exception {\n    IcebergSourceSplit split = splitSerializer.deserialize(splitInfo.version(), splitInfo.split());\n    try (DataIterator<RowData> iterator = rowDataReaderFunction.createDataIterator(split)) {\n      iterator.forEachRemaining(rowData -> extract(rowData, out));\n    } catch (Exception e) {\n      LOG.warn(\"Exception processing split {} at {}\", split, ctx.timestamp(), e);\n      ctx.output(DeleteOrphanFiles.ERROR_STREAM, e);\n      errorCounter.inc();\n    }\n  }\n\n  @Override\n  public void close() throws Exception {\n    super.close();\n    tableLoader.close();\n  }\n\n  /**\n   * Extracts the desired data from the given RowData.\n   *\n   * @param rowData the RowData from which to extract\n   * @param out the Collector to which to output the extracted data\n   */\n  abstract void extract(RowData rowData, Collector<R> out);","sourceCodeStart":83,"sourceCodeEnd":119,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/TableReader.java#L83-L119","documentation":"TableReader processes metadata-table splits (e.g. from the files/manifests metadata table) to feed maintenance actions like orphan-file detection. If deserializing the split or reading its rows fails, the exception is logged with the split and timestamp, forwarded to the ERROR_STREAM side output, and counted instead of failing the job.","triggerScenarios":"processElement receives a SplitInfo whose IcebergSourceSplit cannot be deserialized or whose data iterator throws (corrupt/missing data files, schema mismatch, IO errors) while a DeleteOrphanFiles/ExpireSnapshots pipeline reads the corresponding metadata table.","commonSituations":"Files deleted or corrupted between planning and reading (e.g. concurrent expireSnapshots); schema evolution between job versions and stored split state; transient S3/HDFS read failures; checkpoint-restored splits referencing removed snapshots.","solutions":["Inspect the ERROR_STREAM side output and errorCounter to identify which splits failed and why","Re-run the maintenance action after fixing the underlying IO/snapshot issue so missed files are re-scanned","Avoid running expireSnapshots concurrently with orphan-file detection reading the same table","Check that the job was restored with compatible schema/table state; restart from a fresh savepoint if not"],"exampleFix":"// after\nDataStream<RowData> errorStream = result.getSideOutput(DeleteOrphanFiles.ERROR_STREAM);\nerrorStream.addSink(e -> LOG.error(\"maintenance read failed\", e));","handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"result.getSideOutput(DeleteOrphanFiles.ERROR_STREAM)\n    .addSink(e -> log.error(\"Split processing failed, re-run maintenance action\", e));\n// monitor errorCounter and re-run the action once the underlying issue is fixed","preventionTips":["Do not run expireSnapshots/orphan cleanup concurrently with metadata-table readers","Check the ERROR_STREAM side output and error counter in production monitoring","Restore jobs from savepoints created with the same Iceberg version and schema"],"tags":["flink","iceberg","data-reading","error-handling"],"backgroundTag":"file-read-failed","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"}