{"record":{"id":"cd73ffd57aa4dc80","repo":"apache/iceberg","slug":"failed-to-create-iceberg-input-splits-for-table-cd73ff","errorCode":null,"errorMessage":"Failed to create iceberg input splits for table: ","messagePattern":"Failed to create iceberg input splits for table: ","errorType":"exception","errorClass":"UncheckedIOException","httpStatus":null,"severity":"error","filePath":"flink/v2.2/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.2/flink/src/main/java/org/apache/iceberg/flink/source/FlinkSource.java#L277-L313","documentation":"FlinkSource.build throws this UncheckedIOException when estimating input-split count for source parallelism fails: format.createInputSplits(0) raised an IOException while listing/reading table files. The error message embeds the table being scanned. It surfaces during job construction because parallelism auto-computation requires materializing the splits.","triggerScenarios":"Calling FlinkSource.forRowData()...build() (or SQL equivalents) where underlying FileIO cannot read manifest/data file locations — missing/incorrect warehouse path, expired credentials, deleted table location, or Hadoop/S3 misconfiguration.","commonSituations":"Misconfigured catalog or warehouse path, S3 credentials missing/expired on the jobmanager, table dropped between planning and build, or HDFS NameNode unreachable from the Flink cluster.","solutions":["Verify the table exists and its location is reachable with the configured FileIO (list the location manually).","Check credentials/config for the object store or HDFS in the Flink cluster's environment (HADOOP_CONF_DIR, s3 access keys).","Fix the catalog configuration (type, uri, warehouse) used to load the table.","Set an explicit parallelism to skip the split-count-based auto-parallelism path, avoiding createInputSplits at planning time.","Inspect the wrapped IOException cause for the precise filesystem error (not found vs permission)."],"exampleFix":"// before: relies on auto parallelism probing splits\nFlinkSource.forRowData().env(env).table(table).build();\n// after: pin parallelism to avoid split probing, or fix config first\nFlinkSource.forRowData()\n    .env(env)\n    .table(table)\n    .parallelism(4)\n    .build();","handlingStrategy":"try-catch","validationCode":"// preflight: ensure the table location is listable before building the source\ntable.io().newInputFile(table.location()).getLength();","typeGuard":null,"tryCatchPattern":"try {\n  DataStream<RowData> stream = FlinkSource.forRowData().env(env).table(table).build();\n} catch (UncheckedIOException e) {\n  LOG.error(\"Cannot create splits for table {} — check catalog/warehouse/credentials\", e);\n  throw e;\n}","preventionTips":["Verify catalog/warehouse configuration and filesystem credentials on the Flink cluster before submitting jobs.","Confirm the table location is reachable from jobmanager and taskmanagers.","Set explicit source parallelism to avoid planning-time split materialization.","Check the wrapped IOException cause to distinguish not-found vs permission errors."],"tags":["flink","file-io","filesystem","configuration"],"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"}