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
- Ensure the source does not emit multiple schema-change-before signals before the previous checkpoint trigger completes.
- Check whether the prior triggerSchemaChangeBeforeCheckpoint() future failed or hung; resolve upstream errors keeping the phase set.
- 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.
- 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
- Coalesce/throttle DDL events in CDC sources
- Always reset schemaChangePhase in finally/exception paths
- Await completion of one schema-change checkpoint before emitting the next signal
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
- Agent is running (pid file ); stop the agent before write…
- Airtable API rate limit reached, retry
- At least one source plugin must be configured.
- BigtableSourceSplitEnumerator already closed; cannot create…
- BigtableSourceSplitEnumerator closed during client creation
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)