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

  1. Read the boxed source error for the connector-specific cause.
  2. Check the external system (broker, database) availability, credentials, and network reachability.
  3. Validate the source/sink configuration (topic, bootstrap servers, auth settings).
  4. 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

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


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)