risingwavelabs/risingwave · error · SinkError

The primary key column index

Error message

The primary key column index {} is out of bounds in schema {:?}

What it means

The single primary key index supplied to the encoder points past the end of the sink schema's fields, so no column can be resolved. This is an internal consistency violation between the pk indices and the schema passed into EncoderParams.

Solutions

  1. Recompute pk_indices against the exact schema passed in params.schema
  2. Verify the sink's downstream dispatcher did not drop/reorder columns before encoder build
  3. Ensure pk index refers to a column still present in the sink's output schema

Example fix

// before: stale indices
let encoder = BytesEncoder::build(params_with_full_schema, pk_from_relation).await?;
// after
let pk_indices = pk_from_relation.filter(|i| *i < params.schema.len());
Defensive patterns

Strategy: validation

Validate before calling

let pk = pk_indices.expect("pk required")[0];
assert!(pk < schema.len(), "pk index {} out of bounds for schema len {}", pk, schema.len());

Type guard

fn pk_in_schema(pk: usize, schema: &Schema) -> bool { pk < schema.fields().len() }

Prevention

When it happens

Trigger: build() resolves pk_indices[0] via params.schema.fields().get(pk_indices[0]) and it is None — pk index computed against a different/wider schema than the one given to the encoder.

Common situations: Custom sink code passing a pruned or reordered schema while reusing pk indices from the original relation; version drift where downstream column elimination shrank the schema.

Related errors


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

Appendix: source

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

    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 {
    async fn build(params: EncoderParams<'_>, pk_indices: Option<Vec<usize>>) -> Result<Self> {
        match pk_indices {
            // This is being used as a key encoder
            Some(_) => {
                let (pk_index, schema_ref) = ensure_only_one_pk("BYTES", &params, &pk_indices)?;
                if let DataType::Bytea = schema_ref.data_type() {
                    Ok(BytesEncoder::new(params.schema, pk_index))
                } else {

View on GitHub (pinned to 6469eb736d)