risingwavelabs/risingwave · error · ConnectorError

primary key cannot be null

Error message

primary key {} cannot be null

What it means

During snapshot pagination the connector reads the last primary-key value of the previous chunk to use as the lower bound of the next query. If that PK value comes back as SQL NULL, it cannot be used to build the range predicate, so the connector bails. A NULL primary key indicates corrupted upstream metadata or an mis-detected key column.

Solutions

  1. Ensure the MySQL table has a real PRIMARY KEY (or at minimum a NOT NULL UNIQUE key) before enabling CDC.
  2. Verify the CDC table's configured key columns point to the actual PK columns.
  3. Make the key column NOT NULL upstream (`ALTER TABLE t MODIFY col ... NOT NULL`).
  4. Re-create the CDC table so key metadata is re-derived from the current upstream schema.

Example fix

// before: nullable column treated as key
CREATE TABLE t (code INT, name VARCHAR(50)); // no PK
// after
ALTER TABLE t MODIFY code INT NOT NULL;
ALTER TABLE t ADD PRIMARY KEY (code);
Defensive patterns

Strategy: validation

Validate before calling

// Ensure PK columns are NOT NULL upstream before enabling CDC
SELECT COLUMN_NAME, IS_NULLABLE FROM INFORMATION_SCHEMA.COLUMNS
 WHERE TABLE_SCHEMA = ? AND TABLE_NAME = ? AND COLUMN_KEY = 'PRI';
-- every row must have IS_NULLABLE = 'NO'

Try / catch

match snapshot_read(...).await {
    Err(e) if e.to_string().contains("cannot be null") => {
        // halt ingestion and require upstream PK fix
    },
    r => r?,
}

Prevention

When it happens

Trigger: `snapshot_read_inner` (called from `snapshot_read`) reads a chunk where the value of a column serving as the primary key is NULL (`val.is_null()` branch).

Common situations: Configuring the CDC table against a table with no real primary key so a nullable column is (wrongly) treated as PK; metadata describing the wrong key column; upstream schema drift making the key column nullable.

Related errors


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

Appendix: source

Thrown at src/connector/src/source/cdc/external/mysql.rs:763

                            DataType::Float64 => Value::from(value.into_float64().into_inner()),
                            DataType::Varchar => Value::from(String::from(value.into_utf8())),
                            DataType::Date => Value::from(value.into_date().0),
                            DataType::Time => Value::from(value.into_time().0),
                            DataType::Timestamp => Value::from(value.into_timestamp().0),
                            DataType::Decimal => Value::from(value.into_decimal().to_string()),
                            DataType::Timestamptz => {
                                // Convert timestamptz to NaiveDateTime for MySQL TIMESTAMP comparison
                                // MySQL expects NaiveDateTime for TIMESTAMP parameters
                                let ts = value.into_timestamptz();
                                let datetime_utc = ts.to_datetime_utc();
                                let naive_datetime = datetime_utc.naive_utc();
                                Value::from(naive_datetime)
                            }
                            _ => bail!("unsupported primary key data type: {}", ty),
                        };
                        ConnectorResult::Ok((pk.to_lowercase(), val))
                    } else {
                        bail!("primary key {} cannot be null", pk);
                    }
                })
                .try_collect::<_, _, ConnectorError>()?;

            tracing::debug!("snapshot read params: {:?}", &params);
            let rs_stream = sql
                .with(Params::from(params))
                .stream::<mysql_async::Row, _>(&mut conn)
                .await?;

            let row_stream = rs_stream.map(|row| {
                // convert mysql row into OwnedRow
                let mut row = row?;
                mysql_row_to_owned_row_with_strict_pk(&mut row, &self.rw_schema, &self.pk_indices)
                    .map_err(ConnectorError::from)
            });
            pin_mut!(row_stream);
            #[for_await]

View on GitHub (pinned to 6469eb736d)