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
- 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.
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)