risingwavelabs/risingwave · error · SinkError::Config

mqtt sink encode unsupported

Error message

mqtt sink encode unsupported: {:?}

What it means

Raised in MqttSinkWriter::new when building the row encoder: the FORMAT ... ENCODE ... combination requested for the MQTT sink is not supported by the encoder dispatch. The message includes the debug-formatted unsupported format/encode descriptor so the user can pick a supported one.

Solutions

  1. Use a supported ENCODE for the MQTT sink (JSON, or protobuf with a schema registry/subject id configured).
  2. If using PROTOBUF, provide valid schema registry settings or an inline schema so ProtoEncoder can be built.
Defensive patterns

Strategy: validation

When it happens

Trigger: Thrown at src/connector/src/sink/mqtt.rs:275 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/ec3f41889c364b86. Report an issue: GitHub.

Appendix: source

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

                        &format_desc.options,
                        config.topic.as_deref().unwrap_or(name),
                        None,
                    )
                    .await
                    .map_err(|e| SinkError::Config(anyhow!(e)))?;
                    let header = match sid {
                        None => ProtoHeader::None,
                        Some(sid) => ProtoHeader::ConfluentSchemaRegistry(sid),
                    };
                    RowEncoderWrapper::Proto(ProtoEncoder::new(
                        schema.clone(),
                        None,
                        descriptor,
                        header,
                    )?)
                }
                _ => {
                    return Err(SinkError::Config(anyhow!(
                        "mqtt sink encode unsupported: {:?}",
                        format_desc.encode,
                    )));
                }
            },
            _ => {
                return Err(SinkError::Config(anyhow!(
                    "MQTT sink only supports append-only mode"
                )));
            }
        };
        let qos = config.common.qos();

        let (client, mut eventloop) = config
            .common
            .build_client(actor_id, sink_id.as_raw_id())
            .map_err(|e| SinkError::Mqtt(anyhow!(e)))?;

View on GitHub (pinned to 6469eb736d)