risingwavelabs/risingwave · error

The watermark row should only contain 1 datum

Error message

The watermark row should only contain 1 datum

What it means

`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.

Solutions

  1. Ensure the watermark producer emits a single-column row containing only the watermark value.
  2. Align versions of all components handling watermark encoding/decoding.
  3. If reproducible after a clean redeploy, file an internal bug with the fragment that emits the watermark.
Defensive patterns

Strategy: try-catch

Validate before calling

// producer-side invariant before sending watermark row
assert_eq!(row.len(), 1, "watermark row must contain exactly one datum");

Try / catch

// executor side
match decode_watermark_row(row) {
    Ok(v) => v,
    Err(e) if e.to_string().contains("should only contain 1 datum") => {
        // log the malformed watermark and skip / report internal error
        None.into()
    }
}

Prevention

When it happens

Trigger: `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.

Common situations: 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.

Related errors


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

Appendix: source

Thrown at src/stream/src/executor/watermark_filter.rs:382

            Type::GreaterThanOrEqual,
            DataType::Boolean,
            vec![
                InputRefExpression::new(watermark_type.clone(), event_time_col_idx).boxed(),
                LiteralExpression::new(watermark_type, Some(watermark)).boxed(),
            ],
            eval_error_report,
        )
    }

    fn decode_watermark_row(
        watermark_row: Option<OwnedRow>,
    ) -> StreamExecutorResult<Option<ScalarImpl>> {
        match watermark_row {
            Some(row) => {
                if row.len() == 1 {
                    Ok(row[0].clone())
                } else {
                    bail!("The watermark row should only contain 1 datum");
                }
            }
            _ => Ok(None),
        }
    }

    async fn read_global_max_watermark(
        global_watermark_table: &BatchTable<S>,
        read_epoch: u64,
    ) -> StreamExecutorResult<Option<ScalarImpl>> {
        let global_watermark_iter_futures =
            global_watermark_table
                .vnodes()
                .iter_vnodes()
                .map(|vnode| async move {
                    let pk = row::once(vnode.to_datum());
                    let watermark_row: Option<OwnedRow> = global_watermark_table
                        .get_row(pk, HummockReadEpoch::Committed(read_epoch))

View on GitHub (pinned to 6469eb736d)