{"record":{"id":"06d7d27923db7dc1","repo":"apache/iceberg","slug":"failed-to-create-iceberg-input-splits-for-table-06d7d2","errorCode":null,"errorMessage":"Failed to create iceberg input splits for table: \" + table","messagePattern":"Failed to create iceberg input splits for table: \" \\+ table","errorType":"exception","errorClass":"UncheckedIOException","httpStatus":null,"severity":"error","filePath":"flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/source/FlinkSource.java","lineNumber":295,"sourceCode":"\n    public DataStream<RowData> build() {\n      Preconditions.checkNotNull(env, \"StreamExecutionEnvironment should not be null\");\n      FlinkInputFormat format = buildFormat();\n\n      ScanContext context = contextBuilder.build();\n      TypeInformation<RowData> typeInfo =\n          FlinkCompatibilityUtil.toTypeInfo(FlinkSchemaUtil.convert(context.project()));\n\n      if (!context.isStreaming()) {\n        int parallelism =\n            SourceUtil.inferParallelism(\n                readableConfig,\n                context.limit(),\n                () -> {\n                  try {\n                    return format.createInputSplits(0).length;\n                  } catch (IOException e) {\n                    throw new UncheckedIOException(\n                        \"Failed to create iceberg input splits for table: \" + table, e);\n                  }\n                });\n        if (env.getMaxParallelism() > 0) {\n          parallelism = Math.min(parallelism, env.getMaxParallelism());\n        }\n        return env.createInput(format, typeInfo).setParallelism(parallelism);\n      } else {\n        StreamingMonitorFunction function = new StreamingMonitorFunction(tableLoader, context);\n\n        String monitorFunctionName = String.format(\"Iceberg table (%s) monitor\", table);\n        String readerOperatorName = String.format(\"Iceberg table (%s) reader\", table);\n\n        return env.addSource(function, monitorFunctionName)\n            .transform(readerOperatorName, typeInfo, StreamingReaderOperator.factory(format));\n      }\n    }\n  }","sourceCodeStart":277,"sourceCodeEnd":313,"githubUrl":"https://github.com/apache/iceberg/blob/86d9c8fc543e7c56c9f624eb725f76c9baff9570/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/source/FlinkSource.java#L277-L313","documentation":"FlinkSource.build() wraps IOException from FormatModel.createInputSplits(0) — used to count splits when auto-computing parallelism from the number of scan splits — into an UncheckedIOException with the table name. It means the underlying table scan could not produce input splits during job construction.","triggerScenarios":"Building a FlinkSource with scan.parallelism unset (so parallelism is inferred from split count) while createInputSplits throws IOException — e.g. unreadable metadata files, missing/invalid manifest list, or underlying file system access failure during scan planning.","commonSituations":"Table metadata unreadable due to storage permission problems (S3/HDFS credentials); snapshot referenced by the scan no longer exists after concurrent expiry; network interruption to object storage during planning.","solutions":["Check the wrapped IOException (e.getCause()) — it names the file or location that failed to read during split planning.","Verify FileIO credentials/permissions for the table location (S3 access keys, HDFS delegation tokens).","Confirm the target snapshot still exists; if concurrent expiration deleted it, point the scan to a valid snapshot or branch/tag.","Set an explicit scan.parallelism to bypass split-count-based parallelism inference if planning is transiently failing."],"exampleFix":"// before\nnew FlinkSource.ForRowData().tableLoader(loader).build(); // parallelism inferred, triggers split planning\n// after\nnew FlinkSource.ForRowData().tableLoader(loader).project(schema).configureScanParallelism(4).build();","handlingStrategy":"validation","validationCode":"// verify table location readability and snapshot existence before building the source\nTable table = tableLoader.loadTable();\nSnapshot snap = table.currentSnapshot();\nif (snap == null) throw new IllegalStateException(\"table has no current snapshot\");\ntable.io().newInputFile(snap.manifestListLocation()).getLength(); // proactively read","typeGuard":null,"tryCatchPattern":"try {\n  DataStream<RowData> ds = env.fromSource(source, watermark, \"iceberg\");\n} catch (UncheckedIOException e) {\n  if (e.getMessage().contains(\"Failed to create iceberg input splits\")) {\n    LOG.error(\"split planning failed for table; check IO credentials/snapshot\", e.getCause());\n  }\n  throw e;\n}","preventionTips":["Validate FileIO credentials (S3/HDFS) on the submission node before submitting jobs.","Avoid running snapshot expiry concurrently with job startup.","Set explicit scan.parallelism to skip split-count inference when planning is flaky."],"tags":["flink","source","splits","scan-planning"],"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"}