risingwavelabs/risingwave · error · ConnectorError

primary key column ` ` not found in upstream MySQL primary…

Error message

primary key column `{column_name}` not found in upstream MySQL primary key info

What it means

During CDC snapshot reads, RisingWave builds a key-range comparison strategy for the table's primary key columns. It looks up each PK column in `upstream_mysql_pk_infos` (fetched from MySQL) using case-insensitive name matching; if a requested PK column is absent from that metadata, this error is thrown. It almost always means the upstream PK metadata is stale, incomplete, or was populated for a different table/column set.

Solutions

  1. Re-create or refresh the CDC table so `upstream_mysql_pk_infos` is re-fetched from the current MySQL schema.
  2. Verify the upstream table actually has the expected primary key with `SHOW KEYS FROM <table> / information_schema.KEY_COLUMN_USAGE`.
  3. Check the CDC table definition's PK columns match the upstream table's PK columns exactly.
  4. If the upstream PK was altered, pause ingestion, align the RW table schema, then resume or rebuild the table.

Example fix

// before (stale upstream pk infos for a table whose PK changed in MySQL)
ALTER TABLE orders DROP PRIMARY KEY, ADD PRIMARY KEY (order_uuid);
// RW CDC table still has upstream_mysql_pk_infos for ('id')
// after: rebuild the CDC table so PK info is refreshed
DROP TABLE rw_orders;
CREATE TABLE rw_orders (...) WITH (connector = 'mysql-cdc', table_name = 'orders');
Defensive patterns

Strategy: validation

Validate before calling

// Before creating the CDC table, confirm RW PK columns match upstream MySQL PK info
let upstream_pk_names: Vec<String> = upstream_mysql_pk_infos.iter().map(|(c, _)| c.to_ascii_lowercase()).collect();
for col in rw_pk_columns {
    assert!(upstream_pk_names.contains(&col.to_ascii_lowercase()), "PK column {col} missing from upstream MySQL PK info");
}

Try / catch

match needs_unsigned_i64_compare(col) {
    Ok(flag) => flag,
    Err(_) => { refresh_upstream_pk_infos_and_retry()? }
}

Prevention

When it happens

Trigger: Calling `needs_unsigned_i64_compare` (via `pk_column_unsigned_i64_compare_flags` or `snapshot_read_inner`) with a column name that does not appear in `self.upstream_mysql_pk_infos` — e.g. the table's PK was altered after the CDC table was created, or the upstream info list was populated from a different table or a filtered schema.

Common situations: Altering the MySQL table's primary key while a CDC table exists; renaming a PK column upstream; pointing the CDC source at a table whose PK info was cached from another table; case-only naming differences not handled because upstream info itself is missing the column.

Understand the failure class

Background: 'Could not be found', 'does not exist', 'not found in database': the resource-not-found family when an ID, slug, key, or URI lookup comes back empty — this error's family across 20 libraries.

Related errors


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

Appendix: source

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

        drop(conn);

        Ok(column_infos)
    }

    /// Check whether a column is `BIGINT UNSIGNED`.
    ///
    /// Frontend up-casts narrower unsigned integer types, and non-integer unsigned types
    /// (`FLOAT`/`DOUBLE`/`DECIMAL UNSIGNED`) keep their own comparison semantics. Only
    /// `BIGINT UNSIGNED` can be represented as a negative `i64` in RisingWave and needs
    /// unsigned `u64` comparison/conversion.
    fn needs_unsigned_i64_compare(&self, column_name: &str) -> ConnectorResult<bool> {
        self.upstream_mysql_pk_infos
            .iter()
            .find(|(col_name, _)| col_name.eq_ignore_ascii_case(column_name))
            .map(|(_, col_type)| mysql_type_is_unsigned_bigint(col_type))
            .ok_or_else(|| {
                anyhow!(
                    "primary key column `{column_name}` not found in upstream MySQL primary key info"
                )
                .into()
            })
    }

    /// For each given primary key column (by name), whether it needs unsigned `i64` comparison.
    pub(crate) fn pk_column_unsigned_i64_compare_flags(
        &self,
        pk_names: &[String],
    ) -> ConnectorResult<Vec<bool>> {
        pk_names
            .iter()
            .map(|name| self.needs_unsigned_i64_compare(name))
            .collect()
    }

    /// Convert negative i64 to unsigned u64 based on column type

View on GitHub (pinned to 6469eb736d)