risingwavelabs/risingwave · error · SinkError::Remote
end of response stream
Error message
end of response stream
What it means
The log reader stream terminated normally (returned None) while consume_log_and_sink still expected more responses. RisingWave treats premature end-of-stream of the remote sink response channel as an error because the sink writer must stay alive until the epoch completes.
Solutions
- Check RisingWave and JVM/connector logs around this timestamp for the remote writer closing.
- Verify the downstream sink service was running and reachable for the whole epoch.
- Recover/restart the sink; RisingWave will replay from the last barrier.
- If reproducible, check connector version compatibility between frontend/stream and the remote connector.
Defensive patterns
Strategy: retry
Try / catch
match stream.try_next().await {
Ok(None) => return Err(/* premature EOF: schedule recovery from last barrier */),
Err(e) => return Err(e),
Ok(Some(resp)) => handle(resp),
} Prevention
- Monitor downstream sink service health so it doesn't close streams mid-epoch
- Set adequate gRPC/JNI timeouts so streams aren't closed early
- Alert on sink recovery events correlated with this error
When it happens
Trigger: The remote (JVM) sink writer stream ends (Ok(None) from try_next) while consume_log_and_sink is still polling for responses, e.g. the remote peer closed the stream before all writes were acknowledged.
Common situations: Remote sink process/connector crashed or closed its stream mid-epoch; JNI writer dropped unexpectedly; external sink service stopped while the sink was streaming.
Related errors
- SinkError::Remote(anyhow!(e))
- 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/a8648234feabdbce.
Report an issue: GitHub.
Appendix: source
Thrown at src/connector/src/sink/remote.rs:349
let mut response_err_stream_rx = self.response_stream;
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();
}
View on GitHub (pinned to 6469eb736d)