risingwavelabs/risingwave · error · ConnectorError

Unknown sink connector

Error message

Unknown sink connector: {sink_name}

What it means

Raised by check_sink_allow_alter_on_fly_fields when the sink name is neither the special JdbcSink name nor resolvable via sink_properties::sink_name_to_config_type_name, so no sink type key can be derived for the registry lookup.

Solutions

  1. Check spelling against the registered sink names in sink_properties.
  2. Use an officially supported sink connector name (e.g. 'kafka', 'jdbc', 'snowflake', 'clickhouse', 'elasticsearch').
  3. If adding a new sink, register it in sink_name_to_config_type_name and generate its allow fields.
  4. Align frontend/connector crate versions so both know the sink name.

Example fix

// before
check_sink_allow_alter_on_fly_fields("kafkas", &fields)?;
// after
check_sink_allow_alter_on_fly_fields("kafka", &fields)?;
Defensive patterns

Strategy: validation

Validate before calling

fn is_known_sink(name: &str) -> bool {
    name == JdbcSink::SINK_NAME
        || risingwave_connector::sink::sink_properties::sink_name_to_config_type_name(name).is_some()
}

Try / catch

match check_sink_allow_alter_on_fly_fields(sink, &fields) {
    Ok(()) => {},
    Err(e) => eprintln!("sink '{}' not recognized: {}", sink, e),
}

Prevention

When it happens

Trigger: Calling the sink allow-alter-on-fly validation with an unknown sink name, such as a typo in CREATE SINK connector name or a sink type from a different RisingWave version.

Common situations: Typo in sink connector ('s3' vs 'snowflake' etc.); using a name valid only as a source; sink type renamed/removed across versions.

Understand the failure class

Background: "Not found" and "does not exist" errors: why "Task not found", "No such folder", and "Can't find" fire when a lookup comes back empty — this error's family across 14 libraries.

Related errors


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

Appendix: source

Thrown at src/connector/src/with_options_test.rs:858

}}

/// Checks if all given fields are allowed to be altered on the fly for the specified sink connector.
/// Returns Ok(()) if all fields are allowed, otherwise returns a `ConnectorError`.
pub fn check_sink_allow_alter_on_fly_fields(
    sink_name: &str,
    fields: &[String],
) -> crate::error::ConnectorResult<()> {{
    // TODO(#24846): JDBC sink currently uses `()` as sink config type in `for_all_sinks!`,
    // so it cannot have an isolated key in `SINK_ALLOW_ALTER_ON_FLY_FIELDS`.
    // Reuse the JDBC entry in `CONNECTION_ALLOW_ALTER_ON_FLY_FIELDS` for now.
    // TODO(#24846): remove this special case after JDBC sink has a dedicated config type
    // and allow-alter fields are generated directly into `SINK_ALLOW_ALTER_ON_FLY_FIELDS`.
    let allowed_fields = if sink_name == JdbcSink::SINK_NAME {{
        CONNECTION_ALLOW_ALTER_ON_FLY_FIELDS.get(JdbcSink::SINK_NAME)
    }} else {{
        // Convert sink name to the type name key
        let Some(type_name) = sink_properties::sink_name_to_config_type_name(sink_name) else {{
            return Err(ConnectorError::from(anyhow::anyhow!(
                "Unknown sink connector: {{sink_name}}"
            )));
        }};
        SINK_ALLOW_ALTER_ON_FLY_FIELDS.get(type_name)
    }};
    let Some(allowed_fields) = allowed_fields else {{
        return Err(ConnectorError::from(anyhow::anyhow!(
            "No allow_alter_on_fly fields registered for sink: {{sink_name}}"
        )));
    }};
    for field in fields {{
        if !allowed_fields.contains(field) {{
            return Err(ConnectorError::from(anyhow::anyhow!(
                "Field '{{field}}' is not allowed to be altered on the fly for sink: {{sink_name}}"
            )));
        }}
    }}
    Ok(())

View on GitHub (pinned to 6469eb736d)