risingwavelabs/risingwave · error
current aligning barrier is_checkpoint
Error message
current aligning barrier is_checkpoint: {}, current barrier is_checkpoint {} What it means
Guard in the KV log store reader's check_is_checkpoint: while the stream is in BarrierAligning state, the incoming barrier's is_checkpoint flag differs from the one already being aligned. Two different kinds of barriers cannot be aligned together, so the read fails with this diagnostic message.
Solutions
- check_is_checkpoint detected a mismatch between the checkpoint flag of the barrier being aligned and the current aligning state, i.e., barriers arrived out of order or a barrier was consumed by the wrong stream.
- Check upstream barrier scheduling; restart the actor so barrier alignment state resets.
- If reproducible, report with both flags from the message.
Defensive patterns
Strategy: try-catch
When it happens
Trigger: Thrown at src/stream/src/common/log_store_impl/kv_log_store/serde.rs:606 when the library encounters an invalid state.
Common situations: See trigger scenarios.
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/f1f632de9b6d0581.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/common/log_store_impl/kv_log_store/serde.rs:606
.map(|(vnode, s)| deserialize_stream(vnode, s, serde.clone()).peekable())
.collect(),
row_streams: FuturesUnordered::new(),
not_started_streams: Vec::new(),
stream_state: StreamState::Uninitialized,
metrics,
}
}
fn check_is_checkpoint(&self, is_checkpoint: bool) -> LogStoreResult<()> {
if let StreamState::BarrierAligning {
is_checkpoint: curr_is_checkpoint,
..
} = &self.stream_state
{
if is_checkpoint == *curr_is_checkpoint {
Ok(())
} else {
Err(anyhow!(
"current aligning barrier is_checkpoint: {}, current barrier is_checkpoint {}",
curr_is_checkpoint,
is_checkpoint
))
}
} else {
Ok(())
}
}
#[try_stream(ok = (Epoch, KvLogStoreItem), error = anyhow::Error)]
async fn into_vnode_log_store_item_stream(mut self, chunk_size: usize) {
assert!(chunk_size >= 2, "too small chunk_size: {}", chunk_size);
let mut ops = Vec::with_capacity(chunk_size);
let mut data_chunk_builder =
DataChunkBuilder::new(self.serde.payload_schema.clone(), chunk_size);
let mut progress = HashMap::new();View on GitHub (pinned to 6469eb736d)