risingwavelabs/risingwave · error

auto schema refresh sink must have only one fragment, but…

Error message

auto schema refresh sink must have only one fragment, but got {}

What it means

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.

Solutions

  1. Drop and recreate the sink with auto schema refresh instead of replacing it.
  2. Simplify the sink query so it plans into a single fragment.
  3. Check the sink's fragment layout via fragment view to confirm the count before attempting replace.

Example fix

// before
ALTER SINK complex_sink ...; -- fails: 2 fragments
// after
DROP SINK complex_sink;
CREATE SINK complex_sink AS SELECT col FROM mv; -- simple single-fragment plan
Defensive patterns

Strategy: validation

Validate before calling

// before replace, check fragment count
let frags = meta.get_job_fragments_by_id(sink.id.as_job_id()).await?;
if frags.fragments.len() != 1 {
    return Err(format!("auto schema refresh sink has {} fragments; drop & recreate instead", frags.fragments.len()).into());
}

Try / catch

match replace_job(..).await {
    Err(e) if e.to_string().contains("must have only one fragment") => {
        drop_sink(sink.id).await?; recreate_sink(..).await?;
    }
    other => other?,
}

Prevention

When it happens

Trigger: 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.

Common situations: 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.

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 risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/e61d66e750da0836. Report an issue: GitHub.

Appendix: source

Thrown at src/meta/src/rpc/ddl_controller.rs:1673

                if let Some(variant_column) = first_variant_column(&newly_added_columns) {
                    let sink_names = auto_refresh_schema_sinks
                        .iter()
                        .map(|sink| format!("`{}`", sink.name))
                        .join(", ");
                    return Err(MetaError::invalid_parameter(format!(
                        "cannot add VARIANT column `{}` because sink(s) {} with auto schema refresh do not support VARIANT",
                        variant_column.name_with_hidden(),
                        sink_names,
                    )));
                }
                let mut sinks = Vec::with_capacity(auto_refresh_schema_sinks.len());
                for sink in auto_refresh_schema_sinks {
                    let sink_job_fragments = self
                        .metadata_manager
                        .get_job_fragments_by_id(sink.id.as_job_id())
                        .await?;
                    if sink_job_fragments.fragments.len() != 1 {
                        return Err(anyhow!(
                            "auto schema refresh sink must have only one fragment, but got {}",
                            sink_job_fragments.fragments.len()
                        )
                        .into());
                    }
                    let sink_ctx = sink_job_fragments.ctx;
                    let original_sink_fragment =
                        sink_job_fragments.fragments.into_values().next().unwrap();
                    let (new_sink_fragment, new_schema, new_log_store_table) =
                        rewrite_refresh_schema_sink_fragment(
                            &original_sink_fragment,
                            &sink,
                            &newly_added_columns,
                            &removed_columns,
                            table,
                            fragment_graph.table_fragment_id(),
                            self.env.id_gen_manager(),
                        )?;

View on GitHub (pinned to 6469eb736d)