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
- Recompute pk_indices against the exact schema passed in params.schema
- Verify the sink's downstream dispatcher did not drop/reorder columns before encoder build
- 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
- Recompute pk indices from the exact schema handed to the encoder
- Log schema + pk indices together when debugging sink builds
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
- RisingWave primary key column index
- auto schema refresh sink must have only one fragment, but…
- Can't find data
- Cannot find
- Cannot find
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", ¶ms, &pk_indices)?;
if let DataType::Bytea = schema_ref.data_type() {
Ok(BytesEncoder::new(params.schema, pk_index))
} else {View on GitHub (pinned to 6469eb736d)