risingwavelabs/risingwave · error · anyhow

new response offset does not match the buffer offset

Error message

new response offset {:?} does not match the buffer offset {:?}

What it means

After popping the queue entry for a TruncateOffset response, RisingWave checks that the popped sent_offset equals the persisted_offset the remote side asked to truncate. A mismatch means the remote sink confirmed an offset different from the front of the sent buffer, so ordering is broken.

Solutions

  1. Compare the two offsets in the message to see whether the remote is ahead, behind, or unrelated.
  2. Verify remote connector version matches the RisingWave sink protocol version.
  3. Recover the sink from the last barrier; if it recurs consistently, file a bug with the connector logs.
Defensive patterns

Strategy: try-catch

Try / catch

if sent_offset != persisted_offset {
    // capture both offsets for the bug report, then recover
    log::error!("offset mismatch: remote={persisted_offset:?} local={sent_offset:?}");
    request_recovery_from_last_barrier();
}

Prevention

When it happens

Trigger: Remote sink returns TruncateOffset { offset } where offset != the oldest queued sent_offset, i.e. the remote acknowledged offsets out of order or acknowledged something not in the buffer.

Common situations: Buggy remote connector implementation that truncates out of order; interleaved/duplicated responses after reconnect; version mismatch where offset semantics changed between connector versions.

Related errors


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

Appendix: source

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

        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(|| {
                    anyhow!("get unsent offset {:?} in response", persisted_offset)
                })?;
                if sent_offset != persisted_offset {
                    bail!(
                        "new response offset {:?} does not match the buffer offset {:?}",
                        persisted_offset,
                        sent_offset
                    );
                }

                if let (TruncateOffset::Barrier { .. }, Some(start_time)) =
                    (persisted_offset, start_time)
                {
                    sink_writer_metrics
                        .sink_commit_duration
                        .observe(start_time.elapsed().as_secs_f64());
                }

                log_reader.truncate(persisted_offset)?;
                Ok(())
            }

View on GitHub (pinned to 6469eb736d)