{"record":{"id":"db1ee03235421b8d","repo":"apache/iceberg","slug":"exception-processing-split-at-db1ee0","errorCode":null,"errorMessage":"Exception processing split {} at {}","messagePattern":"Exception processing split (.+?) at (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"flink/v2.2/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.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/TableReader.java#L83-L119","documentation":"TableReader.processElement deserializes an IcebergSourceSplit and iterates its data with a DataIterator, forwarding rows downstream. If reading the split throws, the exception is logged as a warning with the split and timestamp, pushed to the DeleteOrphanFiles ERROR_STREAM, and counted — the operator continues with the next split instead of failing the job, so orphan-file detection proceeds with partial data.","triggerScenarios":"Raised in processElement when splitSerializer.deserialize or rowDataReaderFunction.createDataIterator(split) / iteration throws — file deleted between planning and reading (expired snapshots/orphan cleanup), schema mismatch against the metadata-table schema, or object-store read errors.","commonSituations":"Another job expired snapshots and deleted data/metadata files mid-scan; serialization version mismatch after an Iceberg upgrade of one part of the cluster; transient S3/HDFS read failures; corrupt Parquet/metadata files.","solutions":["Inspect the exception on the DeleteOrphanFiles.ERROR_STREAM side output for the root cause","Ensure no concurrent snapshot-expiration or orphan-cleanup job deletes files while the maintenance job reads them","Align Iceberg versions across the Flink job and cluster to avoid split serialization incompatibilities","Retry the maintenance job; failed splits reduce accuracy of that cycle but do not corrupt state"],"exampleFix":null,"handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"// Read the forwarded exception for diagnosis\nresult.getSideOutput(DeleteOrphanFiles.ERROR_STREAM)\n      .executeAndCollect().forEachRemaining(e -> log.error(\"split read failed\", e));","preventionTips":["Avoid running snapshot expiration/orphan deletion while the reader is active","Keep Iceberg versions consistent between job and cluster to avoid split deserialization mismatches","Retry the maintenance cycle after transient object-store read errors","Verify data files are readable via a small TableScan smoke test"],"tags":["flink","data-reading","split","orphan-files"],"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"}