risingwavelabs/risingwave · error · SinkError::Remote

get unexpected response

Error message

get unexpected response: {:?}

What it means

consume_log_and_sink matches on the response variants it expects (e.g. Commit, TruncateOffset, BeginEpoch/WriteBatch responses). Any other response variant arriving on the stream is a protocol violation and is reported with its Debug representation.

Solutions

  1. Note the Debug-printed response variant; check for a version mismatch between risingwave and the connector/JVM side.
  2. Recover the sink from the last barrier to reset the response stream state.
  3. If a new response variant was added upstream, upgrade RisingWave so the handler recognizes it.
Defensive patterns

Strategy: try-catch

Try / catch

match response {
    Commit{..} | TruncateOffset{..} => handle_expected(response),
    other => {
        log::error!("unexpected sink response: {other:?}");
        // check version skew, then recover from last barrier
        request_recovery_from_last_barrier();
    }
}

Prevention

When it happens

Trigger: A JniSinkWriterStreamResponse arrives that doesn't match the handled variants in the current state of the loop - e.g. a BeginEpoch or Chunk response arriving when the code expects Commit/Truncate handling.

Common situations: Version skew between the stream engine and the embedded JVM connector (new response type added); a corrupted/misordered stream after a reconnect; duplicated epoch-begin responses after recovery.

Related errors


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

Appendix: source

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

                                    Some(sink_writer_stream_response::Response::Commit(
                                        sink_writer_stream_response::CommitResponse {
                                            epoch,
                                            metadata,
                                        },
                                    )),
                            } => {
                                if let Some(metadata) = metadata {
                                    warn!("get unexpected non-empty metadata: {:?}", metadata);
                                }
                                truncate_matched_offset(
                                    &mut sent_offset_queue,
                                    TruncateOffset::Barrier { epoch },
                                    &mut log_reader,
                                    &sink_writer_metrics,
                                )?;
                            }
                            response => {
                                return Err(SinkError::Remote(anyhow!(
                                    "get unexpected response: {:?}",
                                    response
                                )));
                            }
                        }
                    }
                    futures::future::Either::Right(result) => {
                        let (epoch, item): (u64, LogStoreReadItem) = result?;

                        match item {
                            LogStoreReadItem::StreamChunk { chunk, chunk_id } => {
                                let offset = TruncateOffset::Chunk { epoch, chunk_id };
                                if let Some(prev_offset) = &prev_offset {
                                    prev_offset.check_next_offset(offset)?;
                                }
                                let cardinality = chunk.cardinality();
                                sink_writer_metrics
                                    .connector_sink_rows_received

View on GitHub (pinned to 6469eb736d)