risingwavelabs/risingwave · error · SinkError

KEY ENCODE expects only one primary key, but got

Error message

KEY ENCODE {} expects only one primary key, but got {}

What it means

Key encoders in this formatter support encoding only a single-column primary key. The builder throws when pk_indices contains more than one index, because composite primary keys cannot be represented by the selected key encode format.

Solutions

  1. Restructure so the key is a single column (e.g. concat columns into one varchar key)
  2. Use a key encode format that supports composite keys, or add a surrogate single-column PK
  3. Select only one key column in the sink definition

Example fix

// before
CREATE TABLE t (a INT, b INT, PRIMARY KEY (a, b));
CREATE SINK s FROM t WITH (key_encode = 'text');
// after
CREATE SINK s AS SELECT concat(a::varchar, ':', b::varchar) AS k, * FROM t WITH (key_encode = 'text');
Defensive patterns

Strategy: validation

Validate before calling

if let Some(pk) = &sink_pk_indices, pk.len() != 1 {
    return Err(format!("key encode requires exactly 1 pk column, got {}", pk.len()));
}

Type guard

fn is_single_column_pk(pk: &Option<Vec<usize>>) -> bool { matches!(pk, Some(v) if v.len() == 1) }

Prevention

When it happens

Trigger: build() called with pk_indices.len() > 1 — e.g. a sink with KEY ENCODE BYTES/TEXT on a relation with a composite PRIMARY KEY (a, b).

Common situations: Sinking a table declared with a multi-column primary key; joining MVs and inheriting multiple key columns; expecting the encoder to concatenate key columns (it will not).

Understand the failure class

Background: "Must be a positive integer", "Invalid value", "Unsupported": the invalid-argument-value error family, when a library rejects the value you pass — this error's family across 35 libraries.

Related errors


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

Appendix: source

Thrown at src/connector/src/sink/formatter/mod.rs:199

            Some(sid) => ProtoHeader::ConfluentSchemaRegistry(sid),
        };
        ProtoEncoder::new(b.schema, None, descriptor, header)
    }
}

fn ensure_only_one_pk<'a>(
    data_type_name: &'a str,
    params: &'a EncoderParams<'_>,
    pk_indices: &'a Option<Vec<usize>>,
) -> Result<(usize, &'a Field)> {
    let Some(pk_indices) = pk_indices else {
        return Err(SinkError::Config(anyhow!(
            "{}Encoder requires primary key columns to be specified",
            data_type_name
        )));
    };
    if pk_indices.len() != 1 {
        return Err(SinkError::Config(anyhow!(
            "KEY ENCODE {} expects only one primary key, but got {}",
            data_type_name,
            pk_indices.len(),
        )));
    }

    let schema_ref = params.schema.fields().get(pk_indices[0]).ok_or_else(|| {
        SinkError::Config(anyhow!(
            "The primary key column index {} is out of bounds in schema {:?}",
            pk_indices[0],
            params.schema
        ))
    })?;

    Ok((pk_indices[0], schema_ref))
}

impl EncoderBuild for BytesEncoder {

View on GitHub (pinned to 6469eb736d)