risingwavelabs/risingwave · error · SinkError::Config

either topic or topic.field must be set

Error message

either topic or topic.field must be set

What it means

Validation in MqttSink::validate: the MQTT topic target must be specified either as a static `topic` option or dynamically via `topic.field`; neither was set, so the writer would have nowhere to publish. Fires at sink creation/ALTER after the append-only and topic-field checks.

Solutions

  1. Set a static 'topic' in the WITH options, or
  2. Specify 'topic.field' referencing an existing column/field path in the sink schema.
  3. Run the sink validation again after fixing the properties.
Defensive patterns

Strategy: validation

When it happens

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

Appendix: source

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

impl Sink for MqttSink {
    type LogSinker = AsyncTruncateLogSinkerOf<MqttSinkWriter>;

    const SINK_NAME: &'static str = MQTT_SINK;

    crate::impl_validate_sink_unknown_fields!();

    async fn validate(&self) -> Result<()> {
        if !self.is_append_only {
            return Err(SinkError::Mqtt(anyhow!(
                "MQTT sink only supports append-only mode"
            )));
        }

        if let Some(field) = &self.config.topic_field {
            let _ = get_topic_field_index_path(&self.schema, field.as_str())?;
        } else if self.config.topic.is_none() {
            return Err(SinkError::Config(anyhow!(
                "either topic or topic.field must be set"
            )));
        }

        let _client = (self.config.common.build_client(0.into(), 0))
            .context("validate mqtt sink error")
            .map_err(SinkError::Mqtt)?;

        Ok(())
    }

    async fn new_log_sinker(&self, writer_param: SinkWriterParam) -> Result<Self::LogSinker> {
        Ok(MqttSinkWriter::new(
            self.config.clone(),
            self.schema.clone(),
            &self.format_desc,
            &self.name,
            writer_param.sink_id,

View on GitHub (pinned to 6469eb736d)