apache/iceberg · warning

Exception planning scan for

Error message

Exception planning scan for {} at {}

What it means

Warning from the Flink maintenance MetadataTablePlanner when planning a scan over a metadata table (e.g. all_manifests/files) fails. The exception is emitted to the DeleteOrphanFiles ERROR_STREAM and counted, so that planning round produces no splits and downstream cleanup may be skipped for safety.

Solutions

  1. Inspect the side-output ERROR_STREAM for the root exception
  2. Re-run the maintenance action; refresh races are typically transient
  3. Check catalog connectivity and metadata file integrity for the table
  4. Reduce concurrency/worker pool size if planning fails from thread exhaustion
Defensive patterns

Strategy: try-catch

Validate before calling

// Ensure the table is loadable and refreshable before planning
catalog.loadTable(tableIdent).refresh();

Try / catch

try { planSplits(table, scanContext); } catch (Exception e) { LOG.warn("Exception planning scan for {}", table, e); errorCounter.inc(); }

Prevention

When it happens

Trigger: processElement calls table.refresh() then FlinkSplitPlanner.planIcebergSourceSplits; any exception during refresh or planning (snapshot expiry racing, corrupt metadata JSON, thread pool/worker pool errors, catalog timeouts) is caught and logged.

Common situations: Table metadata refreshed while snapshots were concurrently expired; catalog unavailable or slow; corrupted or unreadable manifest list files; worker pool misconfiguration.

Understand the failure class

Background: "API request failed": what wrapped HTTP errors from external APIs mean and how to find the real cause — this error's family across 29 libraries.

Related errors


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

Appendix: source

Thrown at flink/v2.3/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)