risingwavelabs/risingwave · error · SinkError::Config

MQTT sink only supports append-only mode

Error message

MQTT sink only supports append-only mode

What it means

Guard in MqttConfig::from_btreemap: the MQTT sink only supports append-only streams; the sink `type` property must be 'append-only' (or force-append-only). It fires when an upsert or dequeue-cdc sink type is declared for MQTT, since MQTT publishing has no notion of delete/update keys.

Solutions

  1. Create the MQTT sink with type = 'append-only' in the WITH options.
  2. MQTT sink does not support upsert/debezium; if updating semantics are needed, choose a different sink connector.
Defensive patterns

Strategy: validation

When it happens

Trigger: Thrown at src/connector/src/sink/mqtt.rs:158 when the library encounters an invalid state.

Common situations: See trigger scenarios.


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

Appendix: source

Thrown at src/connector/src/sink/mqtt.rs:158

}

// sink write
pub struct MqttSinkWriter {
    pub config: MqttConfig,
    payload_writer: MqttSinkPayloadWriter,
    #[expect(dead_code)]
    schema: Schema,
    encoder: RowEncoderWrapper,
    stopped: Arc<AtomicBool>,
}

/// Basic data types for use with the mqtt interface
impl MqttConfig {
    pub fn from_btreemap(values: BTreeMap<String, String>) -> Result<Self> {
        let config = serde_json::from_value::<MqttConfig>(serde_json::to_value(values).unwrap())
            .map_err(|e| SinkError::Config(anyhow!(e)))?;
        if config.r#type != SINK_TYPE_APPEND_ONLY {
            Err(SinkError::Config(anyhow!(
                "MQTT sink only supports append-only mode"
            )))
        } else {
            Ok(config)
        }
    }
}

impl TryFrom<SinkParam> for MqttSink {
    type Error = SinkError;

    fn try_from(param: SinkParam) -> std::result::Result<Self, Self::Error> {
        let schema = param.schema();
        let config = MqttConfig::from_btreemap(param.properties)?;
        Ok(Self {
            config,
            schema,
            name: param.sink_name,

View on GitHub (pinned to 6469eb736d)