{"record":{"id":"23b5a260b7a78e4b","repo":"risingwavelabs/risingwave","slug":"got-different-schema-change-to-prev-schema-ch","errorCode":null,"errorMessage":"got different schema change {:?} to prev schema change {:?}","messagePattern":"got different schema change (.+?) to prev schema change (.+?)","errorType":"exception","errorClass":"anyhow::Error","httpStatus":null,"severity":"error","filePath":"src/meta/src/manager/sink_coordination/coordinator_worker.rs","lineNumber":823,"sourceCode":"            pending_epochs.entry(epoch).or_default().add_new_request(\n                handle_id,\n                commit_request,\n                self.handle_manager.vnode_bitmap(handle_id),\n            )?;\n            if pending_epochs\n                .first_key_value()\n                .expect(\"non-empty\")\n                .1\n                .aligned()\n            {\n                let (epoch, commit_requests) = pending_epochs.pop_first().expect(\"non-empty\");\n                let mut metadatas = Vec::with_capacity(commit_requests.requests.len());\n                let mut requests = commit_requests.requests.into_iter();\n                let (first_metadata, first_schema_change) = requests.next().expect(\"non-empty\");\n                metadatas.push(first_metadata);\n                for (metadata, schema_change) in requests {\n                    if first_schema_change != schema_change {\n                        return Err(anyhow!(\n                            \"got different schema change {:?} to prev schema change {:?}\",\n                            schema_change,\n                            first_schema_change\n                        ));\n                    }\n                    metadatas.push(metadata);\n                }\n\n                match &mut coordinator {\n                    SinkCommitCoordinator::SinglePhase(coordinator) => {\n                        if !metadatas.is_empty() {\n                            let start_time = Instant::now();\n                            run_future_with_periodic_fn(\n                                coordinator.commit_data(epoch, metadatas).instrument_await(\n                                    Self::commit_span(\"single_phase_commit_data\", sink_id, epoch),\n                                ),\n                                Duration::from_secs(5),\n                                || {","sourceCodeStart":805,"sourceCodeEnd":841,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/manager/sink_coordination/coordinator_worker.rs#L805-L841","documentation":"When collecting a batch of commit requests for an epoch, all handles must report the same schema change (or none). If the first request's schema change differs from a later one, the coordinator cannot produce a single consistent commit and errors out. This guards atomic schema evolution across sink writers.","triggerScenarios":"Different handles in the same commit batch observe different versions of the table schema — e.g. some writers started before a schema change and some after, or a source schema change lands mid-epoch so only part of the writers pick it up.","commonSituations":"Running `ALTER`/schema-change operations on the sink's upstream table while the sink is actively committing; skewed actor deployments after a parallelism change; stale writers replaying old schema metadata.","solutions":["Retry the failed epoch after all writers have restarted with the new schema; the coordinator aborts the commit safely.","Quiesce the sink (pause it) before applying schema changes, then resume so all writers start with the same schema.","Ensure schema-change propagation is atomic — all actors of a fragment should switch schema in the same barrier.","Check whether a parallelism alter mixed writers from old and new plan versions and force a full sink restart."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"// before resuming a sink after schema change, verify all writers report the same schema version\nassert_eq!(writers.iter().map(|w| w.schema_version).collect::<HashSet<_>>().len(), 1);","typeGuard":"fn schemas_consistent(changes: &[Option<SchemaChange>]) -> bool {\n    changes.iter().all(|c| c == &changes[0])\n}","tryCatchPattern":"// on this error, restart the sink so all writers reload the same schema, then retry the epoch\nif err.contains(\"got different schema change\") { restart_sink_with_new_schema(); }","preventionTips":["Pause the sink before applying upstream schema changes; resume after all actors agree on the schema.","Propagate schema changes atomically with a barrier so all writers switch together.","After a parallelism change, force writers onto a single plan version."],"tags":["rust","sink-coordination","schema-change","consistency"],"backgroundTag":"schema-validation-failed","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}