risingwavelabs/risingwave · error · SinkError::Remote

SinkError::Remote(anyhow!(e))

Error message

SinkError::Remote(anyhow!(e))

What it means

This error propagates an error returned by the log reader stream itself (Err(e) branch). The original error e from the remote sink log reader is wrapped in SinkError::Remote and surfaced, so the root cause text is whatever the reader failed with (gRPC/JNI/IO errors).

Solutions

  1. Inspect the wrapped inner message for the true root cause (network vs decode vs JNI).
  2. For transport errors, check network stability between components and retry after recovery.
  3. For decode errors, verify connector/prost message version compatibility.
  4. Enable connector debug logs to capture the failing request.
Defensive patterns

Strategy: try-catch

Try / catch

match log_reader.try_next().await {
    Err(e) => {
        // classify: transport vs decode; retry transport errors, fail fast on decode
        if is_transport(&e) { schedule_retry(); } else { return Err(e.into()); }
    }
    other => handle(other),
}

Prevention

When it happens

Trigger: consume_log_and_sink's try_next() on the log reader returns Err(e); any transport or decoding failure of the sink log stream is re-raised here.

Common situations: gRPC stream to the remote sink writer failed (connection reset, deadline exceeded); JNI layer threw; response message failed to decode.

Related errors


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

Appendix: source

Thrown at src/connector/src/sink/remote.rs:350

        let sink_writer_metrics = self.sink_writer_metrics;

        let (response_tx, mut response_rx) = unbounded_channel();

        let poll_response_stream = async move {
            loop {
                let result = response_err_stream_rx
                    .stream
                    .try_next()
                    .instrument_await("log_sinker_wait_next_response")
                    .await;
                match result {
                    Ok(Some(response)) => {
                        response_tx.send(response).map_err(|err| {
                            SinkError::Remote(anyhow!("unable to send response: {:?}", err.0))
                        })?;
                    }
                    Ok(None) => return Err(SinkError::Remote(anyhow!("end of response stream"))),
                    Err(e) => return Err(SinkError::Remote(anyhow!(e))),
                }
            }
        };

        let poll_consume_log_and_sink = async move {
            fn truncate_matched_offset(
                queue: &mut VecDeque<(TruncateOffset, Option<Instant>)>,
                persisted_offset: TruncateOffset,
                log_reader: &mut impl SinkLogReader,
                sink_writer_metrics: &SinkWriterMetrics,
            ) -> Result<()> {
                while let Some((sent_offset, _)) = queue.front()
                    && sent_offset < &persisted_offset
                {
                    queue.pop_front();
                }

                let (sent_offset, start_time) = queue.pop_front().ok_or_else(|| {

View on GitHub (pinned to 6469eb736d)