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
- Compare the two offsets in the message to see whether the remote is ahead, behind, or unrelated.
- Verify remote connector version matches the RisingWave sink protocol version.
- 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
- Pin the remote connector version to one tested against your RW version
- Add integration tests that assert in-order truncate responses
- Alert on repeated offset mismatches - they indicate connector protocol bugs
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
- get none metadata in commit response for coordinated sink…
- get unexpected response
- get unsent offset in response
- should get AlignInitialEpochResponse but get
- should get Commit response but get
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)