{"record":{"id":"c2df8a06593bd97a","repo":"risingwavelabs/risingwave","slug":"current-epoch-does-not-match-with-decoded-epoch","errorCode":null,"errorMessage":"current epoch {} does not match with decoded epoch {}","messagePattern":"current epoch (.+?) does not match with decoded epoch (.+?)","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/common/log_store_impl/kv_log_store/serde.rs","lineNumber":1080,"sourceCode":"                            epoch: decoded_epoch,\n                            size: read_size,\n                            op: AlignedLogStoreOp::Barrier {\n                                vnodes: Arc::new(aligned_vnodes.finish()),\n                                is_checkpoint,\n                            },\n                        }));\n                    } else {\n                        match &mut self.stream_state {\n                            StreamState::BarrierAligning {\n                                aligned_vnodes,\n                                read_size,\n                                curr_epoch,\n                                is_checkpoint: current_is_checkpoint,\n                            } => {\n                                aligned_vnodes.set(vnode.to_index(), true);\n                                *read_size += size;\n                                if curr_epoch != &decoded_epoch {\n                                    return Err(anyhow!(\n                                        \"current epoch {} does not match with decoded epoch {}\",\n                                        curr_epoch,\n                                        decoded_epoch\n                                    ));\n                                }\n                                if current_is_checkpoint != &is_checkpoint {\n                                    return Err(anyhow!(\n                                        \"current is_checkpoint {} does not match with decoded is_checkpoint {}\",\n                                        current_is_checkpoint,\n                                        is_checkpoint\n                                    ));\n                                }\n                            }\n                            other => {\n                                let mut aligned_vnodes =\n                                    BitmapBuilder::zeroed(self.serde.vnodes().len());\n                                aligned_vnodes.set(vnode.to_index(), true);\n                                *other = StreamState::BarrierAligning {","sourceCodeStart":1062,"sourceCodeEnd":1098,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/common/log_store_impl/kv_log_store/serde.rs#L1062-L1098","documentation":"Thrown during log store replay/deserialization: while scanning buffered rows, the in-memory current tracking state (`curr_epoch` in `BarrierAligning`/`AllConsumingRow`) does not match the `decoded_epoch` of the row just decoded from the log store. The replay logic must see rows grouped by epoch; a decoded row from a different epoch in the middle of processing means the persisted data or the state machine is inconsistent.","triggerScenarios":"During replay of the KV log store (`replay`/`next_barrier` style iteration), when a decoded row's epoch differs from `curr_epoch` tracked in `StreamState` for the currently aligned vnode set.","commonSituations":"Corrupted or partially written log store rows (e.g. crash between writes), manual tampering/copying of the metadata KV store, replaying a log written by a different RisingWave version with different row encoding, or mixing data from an old cluster into a new one.","solutions":["Check whether the state table/log store backend (e.g. the metadata KV store) was corrupted or manually modified; restore from a clean snapshot/backup.","Verify the replaying cluster uses the same version that wrote the log data (encoding/serialization compatibility).","If it occurs after crash recovery, retry recovery from the latest Hummock checkpoint; if reproducible, report with the persisted key range involved.","Inspect decoded row bytes at the failing point for truncated or misaligned payloads."],"exampleFix":"// before: mixing old-version log data into new cluster\nrw meta start --state-store hummock+... (pointing at foreign data)\n// after: restore only data produced by the same cluster and version\nrw meta start --state-store hummock+s3://my-bucket (original cluster data)\n","handlingStrategy":"try-catch","validationCode":"// Before recovery, sanity-check the metadata KV store range for the stream\n// (e.g. iterate rows and confirm epochs are non-decreasing and contiguous)\nfor (key, val) in log_store.scan_range(start..end) {\n    let row = KvLogStoreRow::parse(key, val)?;\n    // compare with previously decoded epoch; abort early on mismatch\n}","typeGuard":"fn decoded_epoch_matches(state_epoch: u64, row: &KvLogStoreRow) -> bool {\n    state_epoch == row.epoch()\n}","tryCatchPattern":"// Rust\nmatch log_store.replay(...).await {\n    Err(e) if e.to_string().contains(\"does not match with decoded epoch\") => {\n        // treat as corrupted log: restore snapshot and retry recovery from last checkpoint\n    }\n    r => r?,\n}","preventionTips":["Use durable, atomic writes for the metadata KV store; avoid manual edits or partial copies.","Keep the RisingWave version consistent between data-writing and replaying clusters.","Take regular Hummock checkpoints so recovery can fall back to a clean epoch."],"tags":["streaming","replay","epoch","deserialization"],"backgroundTag":"internal-invariant-violation","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"}