{"record":{"id":"80a0e51e96500deb","repo":"risingwavelabs/risingwave","slug":"new-response-offset-does-not-match-the-buffer","errorCode":null,"errorMessage":"new response offset {:?} does not match the buffer offset {:?}","messagePattern":"new response offset (.+?) does not match the buffer offset (.+?)","errorType":"exception","errorClass":"anyhow","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/remote.rs","lineNumber":372,"sourceCode":"\n        let poll_consume_log_and_sink = async move {\n            fn truncate_matched_offset(\n                queue: &mut VecDeque<(TruncateOffset, Option<Instant>)>,\n                persisted_offset: TruncateOffset,\n                log_reader: &mut impl SinkLogReader,\n                sink_writer_metrics: &SinkWriterMetrics,\n            ) -> Result<()> {\n                while let Some((sent_offset, _)) = queue.front()\n                    && sent_offset < &persisted_offset\n                {\n                    queue.pop_front();\n                }\n\n                let (sent_offset, start_time) = queue.pop_front().ok_or_else(|| {\n                    anyhow!(\"get unsent offset {:?} in response\", persisted_offset)\n                })?;\n                if sent_offset != persisted_offset {\n                    bail!(\n                        \"new response offset {:?} does not match the buffer offset {:?}\",\n                        persisted_offset,\n                        sent_offset\n                    );\n                }\n\n                if let (TruncateOffset::Barrier { .. }, Some(start_time)) =\n                    (persisted_offset, start_time)\n                {\n                    sink_writer_metrics\n                        .sink_commit_duration\n                        .observe(start_time.elapsed().as_secs_f64());\n                }\n\n                log_reader.truncate(persisted_offset)?;\n                Ok(())\n            }\n","sourceCodeStart":354,"sourceCodeEnd":390,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/remote.rs#L354-L390","documentation":"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.","triggerScenarios":"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.","commonSituations":"Buggy remote connector implementation that truncates out of order; interleaved/duplicated responses after reconnect; version mismatch where offset semantics changed between connector versions.","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."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"if sent_offset != persisted_offset {\n    // capture both offsets for the bug report, then recover\n    log::error!(\"offset mismatch: remote={persisted_offset:?} local={sent_offset:?}\");\n    request_recovery_from_last_barrier();\n}","preventionTips":["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"],"tags":["sink","offset","ordering","protocol"],"backgroundTag":"offset-mismatch","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}