risingwavelabs/risingwave · error · SinkError

Remote sink error

Error message

Remote sink error: {0}

What it means

Variant of `SinkError` used by the remote sink RPC layer: errors returned when the frontend/connector node talks to a remote sink worker over gRPC, or errors relayed from remote sink execution. The anyhow cause is preserved as source and backtrace.

Solutions

  1. Check the connector node is running and reachable; inspect its logs for the original error
  2. Restart/repair the failing connector or compute node; verify cluster health via `risingwave ctl` or meta dashboard
  3. Confirm component versions match after an upgrade (rolling-upgrade mismatch produces remote errors)
Defensive patterns

Strategy: retry

Validate before calling

// health-check the connector node before running remote sinks
grpc_health_probe -addr <connector-node>:50051 || echo "connector node unreachable"

Try / catch

match err {
    SinkError::Remote(e) if is_transient(&e) => retry_with_backoff(),
    SinkError::Remote(e) => { log::error!("remote sink failed: {e:#}"); return Err(e.into()); }
    e => return Err(e.into()),
}

Prevention

When it happens

Trigger: Remote sink communication fails (gRPC unavailable, remote worker panics/returns error) and the remote sink executor wraps it in SinkError::Remote; also used when delegating sink work to a remote connector node that reports failure.

Common situations: Connector node down or unreachable during distributed sink execution; version mismatch between meta/frontend and connector node; network partition in the cluster.

Related errors


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

Appendix: source

Thrown at src/connector/src/sink/mod.rs:1098

        )*
    };
}

def_sink_impl!();

pub type Result<T> = std::result::Result<T, SinkError>;

#[derive(Error, Debug)]
pub enum SinkError {
    #[error("Kafka error: {0}")]
    Kafka(#[from] rdkafka::error::KafkaError),
    #[error("Kinesis error: {0}")]
    Kinesis(
        #[source]
        #[backtrace]
        anyhow::Error,
    ),
    #[error("Remote sink error: {0}")]
    Remote(
        #[source]
        #[backtrace]
        anyhow::Error,
    ),
    #[error("Encode error: {0}")]
    Encode(String),
    #[error("Avro error: {0}")]
    Avro(#[from] apache_avro::Error),
    #[error("Iceberg error: {0}")]
    Iceberg(
        #[source]
        #[backtrace]
        anyhow::Error,
    ),
    #[error("config error: {0}")]
    Config(
        #[source]

View on GitHub (pinned to 6469eb736d)