{"record":{"id":"607fc22c7581353d","repo":"risingwavelabs/risingwave","slug":"lhs-and-rhs-chunk-cardinality-should-be-the-same","errorCode":null,"errorMessage":"lhs and rhs chunk cardinality should be the same","messagePattern":"lhs and rhs chunk cardinality should be the same","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/row_merge.rs","lineNumber":156,"sourceCode":"            }\n        }\n    }\n\n    fn build_chunk(\n        data_types: &[DataType],\n        lhs_mapping: &ColIndexMapping,\n        rhs_mapping: &ColIndexMapping,\n        lhs_chunk: StreamChunk,\n        rhs_chunk: StreamChunk,\n    ) -> Result<Message, StreamExecutorError> {\n        if !(1..=2).contains(&lhs_chunk.cardinality()) {\n            bail!(\"lhs chunk cardinality should be 1 or 2\");\n        }\n        if !(1..=2).contains(&rhs_chunk.cardinality()) {\n            bail!(\"rhs chunk cardinality should be 1 or 2\");\n        }\n        if lhs_chunk.cardinality() != rhs_chunk.cardinality() {\n            bail!(\"lhs and rhs chunk cardinality should be the same\");\n        }\n        let cardinality = lhs_chunk.cardinality();\n        let mut ops = Vec::with_capacity(cardinality);\n        let mut merged_rows = vec![vec![Datum::None; data_types.len()]; cardinality];\n        for (i, (op, lhs_row)) in lhs_chunk.rows().enumerate() {\n            ops.push(op);\n            for (j, d) in lhs_row.iter().enumerate() {\n                // NOTE(kwannoel): Unnecessary columns will not have a mapping,\n                // for instance extra row count column.\n                // those can be skipped here.\n                if let Some(out_index) = lhs_mapping.try_map(j) {\n                    merged_rows[i][out_index] = d.to_owned_datum();\n                }\n            }\n        }\n\n        for (i, (_, rhs_row)) in rhs_chunk.rows().enumerate() {\n            for (j, d) in rhs_row.iter().enumerate() {","sourceCodeStart":138,"sourceCodeEnd":174,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/row_merge.rs#L138-L174","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","solutions":["Log both chunks' cardinalities per epoch and file an issue — this indicates the two inputs lost lockstep synchronization","Recreate the streaming job to rebuild synchronized state","Inspect upstream actors on both sides for dropped/duplicated chunks (check Hummock/barrier logs)","Upgrade RisingWave; cross-input sync bugs in row_merge are fixed in newer releases"],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// Before merging, verify both sides carry the same row count\nfn cardinalities_match(lhs: &StreamChunk, rhs: &StreamChunk) -> bool {\n    lhs.cardinality() == rhs.cardinality()\n}","typeGuard":"fn mergeable_pair(lhs: &StreamChunk, rhs: &StreamChunk) -> bool {\n    (1..=2).contains(&lhs.cardinality()) && lhs.cardinality() == rhs.cardinality()\n}","tryCatchPattern":"match build_result {\n    Err(e) if e.to_string().contains(\"lhs and rhs chunk cardinality should be the same\") => {\n        // Inputs lost lockstep: recreate the job and investigate barrier ordering\n        investigate_barrier_ordering();\n        recreate_streaming_job(job_id);\n    }\n    other => other?,\n}","preventionTips":["Keep both inputs of the row-merge executor on the same barrier cadence","Monitor barrier alignment across fragments; desync precedes this error","Test cluster failover scenarios for two-input executors in staging","Upgrade RisingWave when cross-input synchronization fixes ship"],"tags":["rust","streaming","internal-invariant","row-merge","sync"],"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"}