risingwavelabs/risingwave · error · StreamExecutorError
Connector error
Error message
Connector error: {0} What it means
The ConnectorError variant wraps a BoxedError (boxed source error) under the message 'Connector error: {0}'. It surfaces failures raised by the connector subsystem (sources/sinks: Kafka, CDC, etc.) when data is read or written by a streaming executor. Via From<ConnectorError> the connector's own error is boxed and reported as the source.
Solutions
- Read the boxed source error for the connector-specific cause.
- Check the external system (broker, database) availability, credentials, and network reachability.
- Validate the source/sink configuration (topic, bootstrap servers, auth settings).
- Retry/recover — connectors typically reconnect automatically; persistent errors need config or external-system fixes.
Example fix
// before CREATE SOURCE s (...) WITH (connector='kafka', properties.bootstrap.servers='wrong-host:9092'); // after CREATE SOURCE s (...) WITH (connector='kafka', properties.bootstrap.servers='broker:9092');
Defensive patterns
Strategy: fallback
Validate before calling
// validate connector config before creating the source
fn validate_kafka_props(props: &KafkaProps) -> Result<(), String> {
if props.bootstrap_servers.is_empty() { return Err("bootstrap.servers required".into()); }
if !props.brokers_reachable() { return Err("brokers unreachable".into()); }
Ok(())
} Type guard
fn is_connector_error(e: &StreamExecutorError) -> bool {
e.variant_name() == "ConnectorError"
} Try / catch
if let Err(e) = poll_result {
if e.variant_name() == "ConnectorError" {
tracing::warn!(source = %e, "connector failure; will retry after backoff");
return Ok(()); // stay retryable instead of crashing the actor
}
return Err(e);
} Prevention
- Pre-validate connector options (hosts, topics, credentials) at DDL time.
- Monitor external system health (brokers, upstream DB replication slots).
- Use retry-friendly connector configs for transient outages.
- Keep connector crate and external system versions compatible.
When it happens
Trigger: Thrown when an executor converting ConnectorError into StreamExecutorError (src/stream/src/executor/error.rs:136) hits a connector failure — e.g. a source poll, message parsing, or sink write fails inside the stream pipeline.
Common situations: Kafka broker unavailability or auth failure; CDC connector losing connection to the upstream database; malformed upstream messages; connector rate limits or quota exhaustion.
Related errors
- additional column is not supported for connector
- additional column type
- connector error
- connector ' ' is not supported
- connector not specified when alter sink
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/13987594c3fadd2c.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/executor/error.rs:98
#[from]
#[backtrace]
RpcError,
),
#[error("Channel closed: {0}")]
ChannelClosed(String),
#[error(transparent)]
ExchangeChannelClosed(
#[from]
#[backtrace]
ExchangeChannelClosed,
),
#[error("Failed to align barrier: expected `{0:?}` but got `{1:?}`")]
AlignBarrier(Box<Barrier>, Box<Barrier>),
#[error("Connector error: {0}")]
ConnectorError(
#[source]
#[backtrace]
BoxedError,
),
#[error(transparent)]
DmlError(
#[from]
#[backtrace]
DmlError,
),
#[error(transparent)]
NotImplemented(#[from] NotImplemented),
#[error(transparent)]
Uncategorized(View on GitHub (pinned to 6469eb736d)