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

  1. Check RisingWave and JVM/connector logs around this timestamp for the remote writer closing.
  2. Verify the downstream sink service was running and reachable for the whole epoch.
  3. Recover/restart the sink; RisingWave will replay from the last barrier.
  4. 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

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


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)