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

  1. 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.
  2. Check upstream barrier scheduling; restart the actor so barrier alignment state resets.
  3. 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)