{"record":{"id":"71c23278f63602c2","repo":"risingwavelabs/risingwave","slug":"mqtt-sink-only-supports-append-only-mode","errorCode":null,"errorMessage":"MQTT sink only supports append-only mode","messagePattern":"MQTT sink only supports append-only mode","errorType":"validation","errorClass":"SinkError::Config","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/mqtt.rs","lineNumber":158,"sourceCode":"}\n\n// sink write\npub struct MqttSinkWriter {\n    pub config: MqttConfig,\n    payload_writer: MqttSinkPayloadWriter,\n    #[expect(dead_code)]\n    schema: Schema,\n    encoder: RowEncoderWrapper,\n    stopped: Arc<AtomicBool>,\n}\n\n/// Basic data types for use with the mqtt interface\nimpl MqttConfig {\n    pub fn from_btreemap(values: BTreeMap<String, String>) -> Result<Self> {\n        let config = serde_json::from_value::<MqttConfig>(serde_json::to_value(values).unwrap())\n            .map_err(|e| SinkError::Config(anyhow!(e)))?;\n        if config.r#type != SINK_TYPE_APPEND_ONLY {\n            Err(SinkError::Config(anyhow!(\n                \"MQTT sink only supports append-only mode\"\n            )))\n        } else {\n            Ok(config)\n        }\n    }\n}\n\nimpl TryFrom<SinkParam> for MqttSink {\n    type Error = SinkError;\n\n    fn try_from(param: SinkParam) -> std::result::Result<Self, Self::Error> {\n        let schema = param.schema();\n        let config = MqttConfig::from_btreemap(param.properties)?;\n        Ok(Self {\n            config,\n            schema,\n            name: param.sink_name,","sourceCodeStart":140,"sourceCodeEnd":176,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/mqtt.rs#L140-L176","documentation":"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.","triggerScenarios":"Thrown at src/connector/src/sink/mqtt.rs:158 when the library encounters an invalid state.","commonSituations":"See trigger scenarios.","solutions":["Create the MQTT sink with type = 'append-only' in the WITH options.","MQTT sink does not support upsert/debezium; if updating semantics are needed, choose a different sink connector."],"exampleFix":null,"handlingStrategy":"validation","validationCode":null,"typeGuard":null,"tryCatchPattern":null,"preventionTips":[],"tags":[],"backgroundTag":null,"analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}