{"record":{"id":"87bdc4c8b9c74f86","repo":"risingwavelabs/risingwave","slug":"the-watermark-row-should-only-contain-1-datum","errorCode":null,"errorMessage":"The watermark row should only contain 1 datum","messagePattern":"The watermark row should only contain 1 datum","errorType":"validation","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/watermark_filter.rs","lineNumber":382,"sourceCode":"            Type::GreaterThanOrEqual,\n            DataType::Boolean,\n            vec![\n                InputRefExpression::new(watermark_type.clone(), event_time_col_idx).boxed(),\n                LiteralExpression::new(watermark_type, Some(watermark)).boxed(),\n            ],\n            eval_error_report,\n        )\n    }\n\n    fn decode_watermark_row(\n        watermark_row: Option<OwnedRow>,\n    ) -> StreamExecutorResult<Option<ScalarImpl>> {\n        match watermark_row {\n            Some(row) => {\n                if row.len() == 1 {\n                    Ok(row[0].clone())\n                } else {\n                    bail!(\"The watermark row should only contain 1 datum\");\n                }\n            }\n            _ => Ok(None),\n        }\n    }\n\n    async fn read_global_max_watermark(\n        global_watermark_table: &BatchTable<S>,\n        read_epoch: u64,\n    ) -> StreamExecutorResult<Option<ScalarImpl>> {\n        let global_watermark_iter_futures =\n            global_watermark_table\n                .vnodes()\n                .iter_vnodes()\n                .map(|vnode| async move {\n                    let pk = row::once(vnode.to_datum());\n                    let watermark_row: Option<OwnedRow> = global_watermark_table\n                        .get_row(pk, HummockReadEpoch::Committed(read_epoch))","sourceCodeStart":364,"sourceCodeEnd":400,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/watermark_filter.rs#L364-L400","documentation":"`WatermarkFilterExecutor::decode_watermark_row` expects the watermark row carried in the barrier to contain exactly one datum (the watermark scalar). Rows of any other length indicate a malformed watermark produced upstream and are rejected.","triggerScenarios":"`decode_watermark_row` receives `Some(row)` where `row.len() != 1`, e.g. a watermark channel that serialized a multi-column row or an empty row instead of a single-column watermark value.","commonSituations":"Version skew between components encoding and decoding watermark rows; a bug in the fragment that constructs watermark rows (e.g. attaching a full row instead of the watermark column); internal proto/schema changes.","solutions":["Ensure the watermark producer emits a single-column row containing only the watermark value.","Align versions of all components handling watermark encoding/decoding.","If reproducible after a clean redeploy, file an internal bug with the fragment that emits the watermark."],"exampleFix":null,"handlingStrategy":"try-catch","validationCode":"// producer-side invariant before sending watermark row\nassert_eq!(row.len(), 1, \"watermark row must contain exactly one datum\");","typeGuard":null,"tryCatchPattern":"// executor side\nmatch decode_watermark_row(row) {\n    Ok(v) => v,\n    Err(e) if e.to_string().contains(\"should only contain 1 datum\") => {\n        // log the malformed watermark and skip / report internal error\n        None.into()\n    }\n}","preventionTips":["Keep watermark row encoding single-column across versions","Add round-trip tests for watermark serialization","Upgrade all components together"],"tags":["streaming","watermark","internal-invariant"],"backgroundTag":"unexpected-response-shape","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"}