{"record":{"id":"5837e41daa456551","repo":"apache/seatunnel","slug":"previous-schema-changes-in-progress-schemachangep","errorCode":null,"errorMessage":"previous schema changes in progress, schemaChangePhase: %s","messagePattern":"previous schema changes in progress, schemaChangePhase: (.+?)","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"error","filePath":"seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SourceFlowLifeCycle.java","lineNumber":294,"sourceCode":"                Thread.sleep(IDLE_SLEEP_MS);\n            } else {\n                if (metricsEnabled) {\n                    sourceReadNs.inc(pollCostNs);\n                }\n                collector.resetEmptyThisPollNext();\n                /*\n                 * The current thread obtain a checkpoint lock in the method {@link\n                 * SourceReader#pollNext(Collector)}. When trigger the checkpoint or savepoint,\n                 * other threads try to obtain the lock in the method {@link\n                 * SourceFlowLifeCycle#triggerBarrier(Barrier)}. When high CPU load, checkpoint\n                 * process may be blocked as long time. So we need sleep to free the CPU.\n                 */\n                Thread.sleep(0L);\n            }\n\n            if (collector.captureSchemaChangeBeforeCheckpointSignal()) {\n                if (schemaChangePhase.get() != null) {\n                    throw new IllegalStateException(\n                            \"previous schema changes in progress, schemaChangePhase: \"\n                                    + schemaChangePhase.get());\n                }\n                schemaChangePhase.set(SchemaChangePhase.createBeforePhase());\n                runningTask.triggerSchemaChangeBeforeCheckpoint().get();\n                log.info(\"triggered schema-change-before checkpoint, stopping collect data\");\n            } else if (collector.captureSchemaChangeAfterCheckpointSignal()) {\n                if (schemaChangePhase.get() != null) {\n                    throw new IllegalStateException(\n                            \"previous schema changes in progress, schemaChangePhase: \"\n                                    + schemaChangePhase.get());\n                }\n                schemaChangePhase.set(SchemaChangePhase.createAfterPhase());\n                runningTask.triggerSchemaChangeAfterCheckpoint().get();\n                log.info(\"triggered schema-change-after checkpoint, stopping collect data\");\n            }\n        } else {\n            if (metricsEnabled) {","sourceCodeStart":276,"sourceCodeEnd":312,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SourceFlowLifeCycle.java#L276-L312","documentation":"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.","triggerScenarios":"collector.captureSchemaChangeBeforeCheckpointSignal() returning true twice without the prior phase being cleared (i.e. triggerSchemaChangeBeforeCheckpoint().get() not completing/clearing the phase before the next signal).","commonSituations":"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.","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."],"exampleFix":"// before\n// source emits two DDL events back-to-back, second signal hits in-progress phase\n// after\n// await previous schema-change checkpoint completion before emitting next DDL signal\nawaitPreviousSchemaChangeCheckpoint();\ncollector.captureSchemaChangeBeforeCheckpointSignal();","handlingStrategy":"retry","validationCode":"if (schemaChangePhase.get() != null) {\n    throw new IllegalStateException(\"schema change already in progress: \" + schemaChangePhase.get());\n}","typeGuard":"boolean canStartSchemaChange(AtomicReference<SchemaChangePhase> phase) {\n    return phase.get() == null;\n}","tryCatchPattern":"try {\n    runningTask.triggerSchemaChangeBeforeCheckpoint().get();\n} catch (Exception e) {\n    schemaChangePhase.set(null); // reset marker so future signals are not blocked\n    throw e;\n}","preventionTips":["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"],"tags":["schema-change","concurrency","source"],"backgroundTag":"invalid-state-transition","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}