risingwavelabs/risingwave · error · ConnectorError

Failed to connect to mqtt broker

Error message

Failed to connect to mqtt broker

What it means

The MQTT source enumerator's `list_splits` waits (polling every 500ms) up to 10 seconds for the background MQTT connection to become established, checking `self.connection_check.is_connected()`. If the broker is still not connected after 10 seconds, `bail!` aborts with "Failed to connect to mqtt broker".

Solutions

  1. Verify the broker host/port in the source WITH options (e.g. `broker = 'mqtt://host:1883'`) and test reachability from the RisingWave host (`nc -vz host 1883`).
  2. Check MQTT credentials and TLS settings (ca.pem/cert/key options) are correct.
  3. Confirm the broker is running and accepting connections (mosquitto status / broker logs).
  4. Restart the source after the broker is reachable; the check runs again on each `list_splits`.
  5. If the broker is merely slow (>10s to accept), reduce load or increase availability of the broker — the timeout is hardcoded.

Example fix

// before
WITH (
  connector = 'mqtt',
  broker = 'mqtt://broker.internal:8883'
)  -- TLS configured nowhere; handshake hangs >10s
// after
WITH (
  connector = 'mqtt',
  broker = 'mqtt://broker.internal:8883',
  username = 'rw_user',
  password = '***',
  ca_certificate = '...'
)
Defensive patterns

Strategy: retry

Validate before calling

// Probe broker reachability and TLS before creating the source
use tokio::net::TcpStream;
let addr = "broker.internal:1883";
TcpStream::connect(addr).await
    .context(format!("MQTT broker unreachable at {addr}; fix broker option or network before creating the source"))?;

Try / catch

// Retry source creation with backoff while the broker is down
for attempt in 0..5 {
    match create_mqtt_source(opts.clone()).await {
        Ok(_) => break,
        Err(e) if e.to_string().contains("Failed to connect to mqtt broker") && attempt < 4 => {
            tokio::time::sleep(Duration::from_secs(15 * (attempt + 1))).await;
        }
        Err(e) => return Err(e),
    }
}

Prevention

When it happens

Trigger: Starting an MQTT source where the broker address is wrong/unreachable, TLS handshake fails, credentials are rejected, or the broker is simply slower than 10 seconds to accept the connection — any case where `connection_check` stays disconnected past the 10s deadline.

Common situations: Misconfigured broker URL or port in WITH options; firewall/security group blocking the MQTT port; broker down or restarting; auth (username/password/TLS cert) mismatch; slow networks or heavily loaded brokers exceeding the fixed 10s timeout.

Understand the failure class

Background: ECONNREFUSED and "connection refused" / "could not connect to server" errors: what they mean and how to fix them — this error's family across 44 libraries.

Related errors


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

Appendix: source

Thrown at src/connector/src/source/mqtt/enumerator/mod.rs:152

        // connection_check always gets created if it doesn't exist
        let connection_check = connection_check.unwrap();

        Ok(Self {
            topic: properties.topic,
            broker: broker_url,
            connection_check,
        })
    }

    async fn list_splits(&mut self) -> ConnectorResult<Vec<MqttSplit>> {
        if !self.connection_check.is_connected() {
            let start = std::time::Instant::now();
            loop {
                if self.connection_check.is_connected() {
                    break;
                };
                if start.elapsed().as_secs() > 10 {
                    bail!("Failed to connect to mqtt broker");
                }

                tokio::time::sleep(std::time::Duration::from_millis(500)).await;
            }
        }
        tracing::debug!("found new splits {} for broker {}", self.topic, self.broker);
        Ok(vec![MqttSplit::new(self.topic.clone())])
    }
}

View on GitHub (pinned to 6469eb736d)