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

  1. Capture the reported persisted_offset and compare against the sink's recent write history.
  2. Restart/recover the sink actor so the queue is rebuilt from the source of truth.
  3. 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

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


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)