{"record":{"id":"354db766baf5dbc4","repo":"apache/iceberg","slug":"failed-to-create-iceberg-input-splits-for-table-354db7","errorCode":null,"errorMessage":"Failed to create iceberg input splits for table: {table}","messagePattern":"Failed to create iceberg input splits for table: (.+?)","errorType":"exception","errorClass":"UncheckedIOException","httpStatus":null,"severity":"error","filePath":"flink/v2.1/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.1/flink/src/main/java/org/apache/iceberg/flink/source/FlinkSource.java#L277-L313","documentation":"FlinkSource.build() probes the number of input splits (via format.createInputSplits(0).length) to derive a suggested source parallelism. Any IOException from split creation is wrapped in an UncheckedIOException with this message, meaning the table's files could not be listed/opened during split planning.","triggerScenarios":"Calling FlinkSource.forRowData()...build() where table.io() fails listing data files or opening FileIO — e.g. missing/invalid warehouse path, deleted table location, expired/insufficient S3/HDFS credentials, or unreachable storage.","commonSituations":"Wrong catalog/warehouse configuration pointing at a nonexistent location; cloud credential expiry on long-running clusters; HDFS NameNode unavailability; table dropped between planning and build.","solutions":["Verify the catalog/warehouse path and table location exist and are accessible","Check storage credentials (S3 AK/SK, IAM, HDFS tokens) are valid and not expired","Confirm the underlying filesystem (HDFS NameNode, S3 endpoint) is reachable","Test TableScan planning outside Flink with the same Hadoop conf to isolate the IO failure"],"exampleFix":"// before\nDataStream<RowData> stream = FlinkSource.forRowData().project(schema).tableLoader(loader).build();\n// after\nTable table = loader.loadTable();\ntable.scan().planTasks(); // fails fast with a clearer IO error if storage is unreachable\nDataStream<RowData> stream = FlinkSource.forRowData().project(schema).tableLoader(loader).build();","handlingStrategy":"try-catch","validationCode":"Table table = loader.loadTable();\ntable.scan().planTasks(); // fail fast with clearer IO error before build()","typeGuard":null,"tryCatchPattern":"try {\n  DataStream<RowData> ds = FlinkSource.forRowData()...build();\n} catch (UncheckedIOException e) {\n  throw new RuntimeException(\"Check storage credentials/path: \" + e.getCause(), e);\n}","preventionTips":["Pre-validate table reachability with a small TableScan before job submission","Keep cloud credentials (S3/IAM) refreshed for long-running clusters","Confirm warehouse path and catalog config point to an existing table"],"tags":["io","storage","flink-source"],"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"}