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
- Pass a StreamingJob::Sink when using replace_sink; use a plain create path for non-sink jobs.
- Fix the calling code so ALTER SINK builds the correct StreamingJob::Sink variant.
- 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
- Only pass replace_sink when the job is genuinely a sink
- Type-check StreamingJob variant before invoking replace paths
- Use dedicated replace/alter APIs for non-sink objects
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
- column not found in the downstream SQL Server table
- `commit_checkpoint_interval` must be greater than 0
- `commit_checkpoint_interval` must be greater than 0
- config error
- connector not specified when alter sink
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)