apache/seatunnel · error · IllegalStateException

previous schema changes in progress, schemaChangePhase

Error message

previous schema changes in progress, schemaChangePhase: %s

What it means

SourceFlowLifeCycle.collect uses schemaChangePhase as an in-progress marker around schema-change-before checkpoints. If a new before-checkpoint signal arrives while schemaChangePhase is still set, it throws IllegalStateException 'previous schema changes in progress, schemaChangePhase: ...'. This enforces that only one schema-change phase runs at a time.

Solutions

  1. Ensure the source does not emit multiple schema-change-before signals before the previous checkpoint trigger completes.
  2. Check whether the prior triggerSchemaChangeBeforeCheckpoint() future failed or hung; resolve upstream errors keeping the phase set.
  3. Inspect the reported schemaChangePhase in the message to identify which phase was never cleared; report as a bug if the engine failed to reset it.
  4. Throttle or coalesce DDL events in the source connector so signals arrive sequentially.

Example fix

// before
// source emits two DDL events back-to-back, second signal hits in-progress phase
// after
// await previous schema-change checkpoint completion before emitting next DDL signal
awaitPreviousSchemaChangeCheckpoint();
collector.captureSchemaChangeBeforeCheckpointSignal();
Defensive patterns

Strategy: retry

Validate before calling

if (schemaChangePhase.get() != null) {
    throw new IllegalStateException("schema change already in progress: " + schemaChangePhase.get());
}

Type guard

boolean canStartSchemaChange(AtomicReference<SchemaChangePhase> phase) {
    return phase.get() == null;
}

Try / catch

try {
    runningTask.triggerSchemaChangeBeforeCheckpoint().get();
} catch (Exception e) {
    schemaChangePhase.set(null); // reset marker so future signals are not blocked
    throw e;
}

Prevention

When it happens

Trigger: collector.captureSchemaChangeBeforeCheckpointSignal() returning true twice without the prior phase being cleared (i.e. triggerSchemaChangeBeforeCheckpoint().get() not completing/clearing the phase before the next signal).

Common situations: Rapid consecutive schema-change events from the source (e.g. DDL storms in CDC), or a prior schema-change trigger that failed/hung leaving the phase marker set.

Understand the failure class

Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/5837e41daa456551. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SourceFlowLifeCycle.java:294

                Thread.sleep(IDLE_SLEEP_MS);
            } else {
                if (metricsEnabled) {
                    sourceReadNs.inc(pollCostNs);
                }
                collector.resetEmptyThisPollNext();
                /*
                 * The current thread obtain a checkpoint lock in the method {@link
                 * SourceReader#pollNext(Collector)}. When trigger the checkpoint or savepoint,
                 * other threads try to obtain the lock in the method {@link
                 * SourceFlowLifeCycle#triggerBarrier(Barrier)}. When high CPU load, checkpoint
                 * process may be blocked as long time. So we need sleep to free the CPU.
                 */
                Thread.sleep(0L);
            }

            if (collector.captureSchemaChangeBeforeCheckpointSignal()) {
                if (schemaChangePhase.get() != null) {
                    throw new IllegalStateException(
                            "previous schema changes in progress, schemaChangePhase: "
                                    + schemaChangePhase.get());
                }
                schemaChangePhase.set(SchemaChangePhase.createBeforePhase());
                runningTask.triggerSchemaChangeBeforeCheckpoint().get();
                log.info("triggered schema-change-before checkpoint, stopping collect data");
            } else if (collector.captureSchemaChangeAfterCheckpointSignal()) {
                if (schemaChangePhase.get() != null) {
                    throw new IllegalStateException(
                            "previous schema changes in progress, schemaChangePhase: "
                                    + schemaChangePhase.get());
                }
                schemaChangePhase.set(SchemaChangePhase.createAfterPhase());
                runningTask.triggerSchemaChangeAfterCheckpoint().get();
                log.info("triggered schema-change-after checkpoint, stopping collect data");
            }
        } else {
            if (metricsEnabled) {

View on GitHub (pinned to cf67b549a7)