risingwavelabs/risingwave · error

lhs and rhs chunk cardinality should be the same

Error message

lhs and rhs chunk cardinality should be the same

What it means

After both chunks individually pass the 1-or-2-row check, `build_chunk` requires equal cardinality on both sides since it merges lhs row i with rhs row i in lockstep. Mismatched cardinalities mean the two inputs lost synchronization and there is no defined pairing, so the executor bails.

Solutions

  1. Log both chunks' cardinalities per epoch and file an issue — this indicates the two inputs lost lockstep synchronization
  2. Recreate the streaming job to rebuild synchronized state
  3. Inspect upstream actors on both sides for dropped/duplicated chunks (check Hummock/barrier logs)
  4. Upgrade RisingWave; cross-input sync bugs in row_merge are fixed in newer releases
Defensive patterns

Strategy: try-catch

Validate before calling

// Before merging, verify both sides carry the same row count
fn cardinalities_match(lhs: &StreamChunk, rhs: &StreamChunk) -> bool {
    lhs.cardinality() == rhs.cardinality()
}

Type guard

fn mergeable_pair(lhs: &StreamChunk, rhs: &StreamChunk) -> bool {
    (1..=2).contains(&lhs.cardinality()) && lhs.cardinality() == rhs.cardinality()
}

Try / catch

match build_result {
    Err(e) if e.to_string().contains("lhs and rhs chunk cardinality should be the same") => {
        // Inputs lost lockstep: recreate the job and investigate barrier ordering
        investigate_barrier_ordering();
        recreate_streaming_job(job_id);
    }
    other => other?,
}

Prevention

When it happens

Trigger: Internal: lhs and rhs chunks for the same merge step contain different row counts (e.g. lhs has 2 rows, rhs has 1), produced by desynchronized buffering between the two inputs for the same epoch.

Common situations: Epoch/barrier desync between the two upstream fragments; one side emitting extra chunks per epoch; regressions in the buffering that equalizes chunk sizes before merge.

Understand the failure class

Background: "This is a bug, please report it": internal invariant violations, unreachable panics, and SNH errors explained — this error's family across 47 libraries.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/607fc22c7581353d. Report an issue: GitHub.

Appendix: source

Thrown at src/stream/src/executor/row_merge.rs:156

            }
        }
    }

    fn build_chunk(
        data_types: &[DataType],
        lhs_mapping: &ColIndexMapping,
        rhs_mapping: &ColIndexMapping,
        lhs_chunk: StreamChunk,
        rhs_chunk: StreamChunk,
    ) -> Result<Message, StreamExecutorError> {
        if !(1..=2).contains(&lhs_chunk.cardinality()) {
            bail!("lhs chunk cardinality should be 1 or 2");
        }
        if !(1..=2).contains(&rhs_chunk.cardinality()) {
            bail!("rhs chunk cardinality should be 1 or 2");
        }
        if lhs_chunk.cardinality() != rhs_chunk.cardinality() {
            bail!("lhs and rhs chunk cardinality should be the same");
        }
        let cardinality = lhs_chunk.cardinality();
        let mut ops = Vec::with_capacity(cardinality);
        let mut merged_rows = vec![vec![Datum::None; data_types.len()]; cardinality];
        for (i, (op, lhs_row)) in lhs_chunk.rows().enumerate() {
            ops.push(op);
            for (j, d) in lhs_row.iter().enumerate() {
                // NOTE(kwannoel): Unnecessary columns will not have a mapping,
                // for instance extra row count column.
                // those can be skipped here.
                if let Some(out_index) = lhs_mapping.try_map(j) {
                    merged_rows[i][out_index] = d.to_owned_datum();
                }
            }
        }

        for (i, (_, rhs_row)) in rhs_chunk.rows().enumerate() {
            for (j, d) in rhs_row.iter().enumerate() {

View on GitHub (pinned to 6469eb736d)