risingwavelabs/risingwave · error · anyhow
get unsent offset in response
Error message
get unsent offset {:?} in response What it means
truncate_matched_offset pops the internal queue of already-sent (offset, timestamp) entries to match the persisted offset reported by a TruncateOffset response. If the queue is empty when it still needs an unsent offset, the bookkeeping invariant is broken and this error is raised.
Solutions
- Capture the reported persisted_offset and compare against the sink's recent write history.
- Restart/recover the sink actor so the queue is rebuilt from the source of truth.
- If reproducible, report with the remote connector name - the connector likely acknowledged offsets it never received.
Defensive patterns
Strategy: try-catch
Try / catch
if let Err(e) = truncate_matched_offset(&mut queue, offset).await {
log::warn!("truncate bookkeeping broken: {e}; resetting from last barrier");
request_recovery_from_last_barrier();
} Prevention
- Ensure the remote connector only acknowledges offsets it actually received
- Rebuild sent-offset queue state deterministically on recovery
- Add debug logging of truncate requests/responses to catch duplicates
When it happens
Trigger: A TruncateOffset { offset } response arrives whose offset has no remaining entry in the sent-offset queue - the queue was drained by earlier truncations or never contained this offset.
Common situations: Duplicate or out-of-order truncate responses from the remote sink; recovery replay that loses in-memory queue state; a bug in the remote connector acknowledging offsets it was never sent.
Understand the failure class
Background: Record Not Found Errors: "not found", RecordNotFound, and "was not found" — what they mean and how to fix them — this error's family across 28 libraries.
Related errors
- Can't create deltalake sink write result from empty data!
- new response offset does not match the buffer offset
- All valid CDC connectors should have returned by now
- ALTER SINK_RATE_LIMIT is not for sink into table
- ambiguous auth: multiple auth options provided; remove one…
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/8d56e674ffb399b2.
Report an issue: GitHub.
Appendix: source
Thrown at src/connector/src/sink/remote.rs:369
}
}
};
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)?;View on GitHub (pinned to 6469eb736d)