risingwavelabs/risingwave · error

replace sink requires a sink job

Error message

replace sink requires a sink job

What it means

The replace_sink path updates an existing sink by resubmitting a streaming job, but that job must actually be a sink. create_streaming_job checks that when replace_sink is Some, the streaming_job matches StreamingJob::Sink; otherwise the operation is invalid and bails.

Solutions

  1. Pass a StreamingJob::Sink when using replace_sink; use a plain create path for non-sink jobs.
  2. Fix the calling code so ALTER SINK builds the correct StreamingJob::Sink variant.
  3. If you meant to replace a table/index job, use the dedicated ALTER/replace paths for those objects.

Example fix

// before
controller.create_streaming_job(StreamingJob::Table(..), None, Some(sink_id), ..).await?;
// after
controller.create_streaming_job(StreamingJob::Sink(sink, ..), None, Some(sink_id), ..).await?;
Defensive patterns

Strategy: validation

Validate before calling

// caller-side check
if replace_sink.is_some() && !matches!(streaming_job, StreamingJob::Sink(..)) {
    return Err("replace_sink requires StreamingJob::Sink".into());
}

Type guard

fn is_sink_job(job: &StreamingJob) -> bool {
    matches!(job, StreamingJob::Sink(..))
}

Try / catch

match create_streaming_job(job, ..).await {
    Err(e) if e.to_string().contains("replace sink requires a sink job") => fix_job_type(),
    other => other?,
}

Prevention

When it happens

Trigger: Calling the DDL controller's create_streaming_job with a non-sink StreamingJob (table, index, MV) while passing replace_sink = Some(old_sink_id).

Common situations: Client/driver bugs mapping ALTER SINK to the wrong streaming job type; custom tooling reusing the create API for replacement; refactors that change StreamingJob variants without updating replace logic.

Understand the failure class

Background: "Must be a positive integer", "Invalid value", "Unsupported": the invalid-argument-value error family, when a library rejects the value you pass — this error's family across 35 libraries.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/97956a9c13d804a2. Report an issue: GitHub.

Appendix: source

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

    }

    /// For [`CreateType::Foreground`], the function will only return after backfilling finishes
    /// ([`crate::manager::MetadataManager::wait_streaming_job_finished`]).
    #[await_tree::instrument(boxed, "create_streaming_job({streaming_job})")]
    pub async fn create_streaming_job(
        &self,
        mut streaming_job: StreamingJob,
        fragment_graph: StreamFragmentGraphProto,
        dependencies: HashSet<ObjectId>,
        resource_type: streaming_job_resource_type::ResourceType,
        if_not_exists: bool,
        refresh_interval_sec: Option<u64>,
        replace_sink: Option<SinkId>,
        since_timestamp_epoch: Option<u64>,
    ) -> MetaResult<NotificationVersion> {
        let replace_sink_info = if let Some(old_sink_id) = replace_sink {
            let StreamingJob::Sink(sink, _) = &streaming_job else {
                bail!("replace sink requires a sink job")
            };
            if sink.target_table.is_some() {
                bail_not_implemented!("replace sink into table")
            }

            Some(old_sink_id)
        } else {
            if let StreamingJob::Sink(sink, _) = &streaming_job
                && let Some(target_table) = sink.target_table
            {
                self.validate_table_for_sink(target_table).await?;
            }
            None
        };
        self.validate_serverless_backfill_enabled(&resource_type)?;
        let ctx = StreamContext::from_protobuf(fragment_graph.get_ctx().unwrap());
        let adaptive_parallelism_strategy =
            (!fragment_graph.adaptive_parallelism_strategy.is_empty()).then(|| {

View on GitHub (pinned to 6469eb736d)