{"record":{"id":"e61d66e750da0836","repo":"risingwavelabs/risingwave","slug":"auto-schema-refresh-sink-must-have-only-one-fragme","errorCode":null,"errorMessage":"auto schema refresh sink must have only one fragment, but got {}","messagePattern":"auto schema refresh sink must have only one fragment, but got (.+?)","errorType":"validation","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/meta/src/rpc/ddl_controller.rs","lineNumber":1673,"sourceCode":"                if let Some(variant_column) = first_variant_column(&newly_added_columns) {\n                    let sink_names = auto_refresh_schema_sinks\n                        .iter()\n                        .map(|sink| format!(\"`{}`\", sink.name))\n                        .join(\", \");\n                    return Err(MetaError::invalid_parameter(format!(\n                        \"cannot add VARIANT column `{}` because sink(s) {} with auto schema refresh do not support VARIANT\",\n                        variant_column.name_with_hidden(),\n                        sink_names,\n                    )));\n                }\n                let mut sinks = Vec::with_capacity(auto_refresh_schema_sinks.len());\n                for sink in auto_refresh_schema_sinks {\n                    let sink_job_fragments = self\n                        .metadata_manager\n                        .get_job_fragments_by_id(sink.id.as_job_id())\n                        .await?;\n                    if sink_job_fragments.fragments.len() != 1 {\n                        return Err(anyhow!(\n                            \"auto schema refresh sink must have only one fragment, but got {}\",\n                            sink_job_fragments.fragments.len()\n                        )\n                        .into());\n                    }\n                    let sink_ctx = sink_job_fragments.ctx;\n                    let original_sink_fragment =\n                        sink_job_fragments.fragments.into_values().next().unwrap();\n                    let (new_sink_fragment, new_schema, new_log_store_table) =\n                        rewrite_refresh_schema_sink_fragment(\n                            &original_sink_fragment,\n                            &sink,\n                            &newly_added_columns,\n                            &removed_columns,\n                            table,\n                            fragment_graph.table_fragment_id(),\n                            self.env.id_gen_manager(),\n                        )?;","sourceCodeStart":1655,"sourceCodeEnd":1691,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/rpc/ddl_controller.rs#L1655-L1691","documentation":"During replace_job, sinks with auto schema refresh enabled must map to exactly one fragment so the replacement can rewrite it atomically. If the sink job's fragment count differs from 1, the replace is rejected with the actual fragment count in the message.","triggerScenarios":"Replacing (ALTER) a job whose auto-refresh-schema sink has more than one fragment — typically a sink that was planned with extra fragments (aggregation, distribution) rather than a single simple fragment.","commonSituations":"Older sinks created before single-fragment guarantees; sinks with complex pipelines that compiled to multiple fragments; clusters where the sink plan changed between creation and replacement.","solutions":["Drop and recreate the sink with auto schema refresh instead of replacing it.","Simplify the sink query so it plans into a single fragment.","Check the sink's fragment layout via fragment view to confirm the count before attempting replace."],"exampleFix":"// before\nALTER SINK complex_sink ...; -- fails: 2 fragments\n// after\nDROP SINK complex_sink;\nCREATE SINK complex_sink AS SELECT col FROM mv; -- simple single-fragment plan","handlingStrategy":"validation","validationCode":"// before replace, check fragment count\nlet frags = meta.get_job_fragments_by_id(sink.id.as_job_id()).await?;\nif frags.fragments.len() != 1 {\n    return Err(format!(\"auto schema refresh sink has {} fragments; drop & recreate instead\", frags.fragments.len()).into());\n}","typeGuard":null,"tryCatchPattern":"match replace_job(..).await {\n    Err(e) if e.to_string().contains(\"must have only one fragment\") => {\n        drop_sink(sink.id).await?; recreate_sink(..).await?;\n    }\n    other => other?,\n}","preventionTips":["Keep auto-schema-refresh sink queries simple so they plan to one fragment","Prefer drop+recreate for complex sinks","Check fragment counts before attempting replace"],"tags":["sink","schema","replace"],"backgroundTag":"invalid-state-transition","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"}