risingwavelabs/risingwave · error · SinkError::DynamoDb

RisingWave primary key column index

Error message

RisingWave primary key column index {} is out of range

What it means

rw_pk_names maps each primary-key column index to its field name using the sink schema. If a pk index is >= the number of schema fields, .get() returns None and this error is thrown. It indicates the pk_indices no longer line up with the schema supplied to the sink.

Solutions

  1. Recreate the sink against the current relation schema so pk_indices match the schema.
  2. Verify the upstream relation's primary key still exists and note its new column positions.
  3. Report as a bug if it reproduces with an unmodified relation, since it indicates an internal invariant violation.
  4. Restore/roll back the upstream schema change that removed or reordered PK columns.
Defensive patterns

Strategy: try-catch

Validate before calling

// sanity-check pk indices against schema length before use
pk_indices.iter().all(|&i| i < schema.fields().len())

Try / catch

// treat as internal invariant violation; log schema + pk_indices and recreate the sink
match sink.validate().await {
  Err(e) if e.to_string().contains("out of range") => recreate_sink_with_current_schema(),
  Err(e) => return Err(e),
  Ok(()) => {},
}

Prevention

When it happens

Trigger: Internal wiring mismatch: pk_indices referencing columns absent from the schema passed to validate (e.g. after schema change on the upstream relation while a sink exists).

Common situations: Upstream table/materialized view altered (columns dropped/reordered) after sink creation; connector-internal invariant violation during sink recovery with a stale schema.

Related errors


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

Appendix: source

Thrown at src/connector/src/sink/dynamodb.rs:407

            return Err(SinkError::DynamoDb(anyhow!("map is not supported yet")));
        }
        DataType::Vector(_) => {
            return Err(SinkError::DynamoDb(anyhow!("vector is not supported yet")));
        }
    };
    Ok(attr)
}

fn rw_pk_names(schema: &Schema, pk_indices: &[usize]) -> Result<Vec<String>> {
    pk_indices
        .iter()
        .map(|pk_idx| {
            schema
                .fields()
                .get(*pk_idx)
                .map(|field| field.name.clone())
                .ok_or_else(|| {
                    SinkError::DynamoDb(anyhow!(
                        "RisingWave primary key column index {} is out of range",
                        pk_idx
                    ))
                })
        })
        .collect()
}

fn dynamodb_key_schema_names(
    table_name: &str,
    key_schema: &[KeySchemaElement],
) -> Result<Vec<String>> {
    if key_schema.is_empty() {
        return Err(SinkError::DynamoDb(anyhow!(
            "table {} key schema is empty",
            table_name
        )));
    }

View on GitHub (pinned to 6469eb736d)