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
- 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`).
- Check MQTT credentials and TLS settings (ca.pem/cert/key options) are correct.
- Confirm the broker is running and accepting connections (mosquitto status / broker logs).
- Restart the source after the broker is reachable; the check runs again on each `list_splits`.
- 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
- Validate broker host/port/network reachability (firewall, security groups) before creating the source.
- Double-check username/password and TLS certificate options against broker config.
- Ensure the broker is running and sized to accept connections within ~10 seconds.
- Monitor broker health and restarts; treat this error as a broker availability signal.
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
- Mqtt error
- SinkError::Mqtt(anyhow!(e))
- all request confluent registry all timeout
- {batch write failure, with context()}
- bigquery insert error
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)