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
- Inspect the wrapped inner message for the true root cause (network vs decode vs JNI).
- For transport errors, check network stability between components and retry after recovery.
- For decode errors, verify connector/prost message version compatibility.
- 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
- Monitor network stability between RW and the remote connector
- Keep prost/connector message versions aligned to avoid decode errors
- Log the inner error immediately - it carries the true root cause
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
- end of response stream
- Cannot find ' ',please set it.
- end of stream
- sink validation failed
- ` ` and ` ` only one can be set
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)