apache/iceberg · warning

Exception planning scan for

Error message

Exception planning scan for {} at {}

What it means

MetadataTablePlanner.processElement plans scans over Iceberg metadata tables (e.g. all_manifests, all_data_files) for the orphan-files flow and emits serialized splits. If planning fails — table refresh, scan planning via FlinkSplitPlanner, or serialization — the exception is logged with the table and timestamp, sent to the DeleteOrphanFiles ERROR_STREAM, and counted; the maintenance cycle continues without splits for that trigger.

Solutions

  1. Inspect the exception on the ERROR_STREAM side output to identify whether refresh or split planning failed
  2. Verify catalog access and that the table metadata is readable from the TaskManager
  3. Increase planning resources (worker pool size, TaskManager memory) for very large metadata tables
  4. Retry the maintenance job — the next trigger re-plans with a refreshed table
Defensive patterns

Strategy: retry

Try / catch

// Same side-output pattern as other maintenance operators
result.getSideOutput(DeleteOrphanFiles.ERROR_STREAM)
      .executeAndCollect().forEachRemaining(e -> log.error("plan failed", e));

Prevention

When it happens

Trigger: Raised in processElement after table.refresh() or FlinkSplitPlanner.planIcebergSourceSplits(table, scanContext, workerPool) throws — invalid scan context, catalog refresh failure, corrupted metadata, or worker-pool errors during split planning.

Common situations: Catalog connectivity or auth failure on refresh; very large tables exhausting planning memory/time; snapshot expired concurrently and removed mid-plan; invalid DeleteOrphanFiles scan configuration (e.g. bad maxCurrentlyReferencedFiles-style options).

Related errors


AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12). Data as JSON: /api/errors/1dc4340307ae1e84. Report an issue: GitHub.

Appendix: source

Thrown at flink/v2.2/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/MetadataTablePlanner.java:101

        ThreadPools.newFixedThreadPool(table.name() + "-table-planner", workerPoolSize);
    this.splitSerializer = new IcebergSourceSplitSerializer(scanContext.caseSensitive());
    this.errorCounter =
        TableMaintenanceMetrics.groupFor(
                getRuntimeContext(), originalTable.name(), taskName, taskIndex)
            .counter(TableMaintenanceMetrics.ERROR_COUNTER);
  }

  @Override
  public void processElement(Trigger trigger, Context ctx, Collector<SplitInfo> out)
      throws Exception {
    try {
      table.refresh();
      for (IcebergSourceSplit split :
          FlinkSplitPlanner.planIcebergSourceSplits(table, scanContext, workerPool)) {
        out.collect(new SplitInfo(splitSerializer.getVersion(), splitSerializer.serialize(split)));
      }
    } catch (Exception e) {
      LOG.warn("Exception planning scan for {} at {}", table, ctx.timestamp(), e);
      ctx.output(DeleteOrphanFiles.ERROR_STREAM, e);
      errorCounter.inc();
    }
  }

  @Override
  public void close() throws Exception {
    super.close();
    tableLoader.close();
    if (workerPool != null) {
      workerPool.shutdown();
    }
  }

  public static class SplitInfo {
    private final int version;
    private final byte[] split;

View on GitHub (pinned to 86d9c8fc54)