risingwavelabs/risingwave · error · SinkError::Config
Primary key column not found in Iceberg schema
Error message
Primary key column {} not found in Iceberg schema What it means
Thrown by resolve_equality_delete_field_ids when a primary key column name from the sink's PRIMARY KEY definition cannot be resolved to a field id in the target Iceberg table schema. Iceberg identifies columns by integer field ids, not names, so an unresolvable name means the sink configuration does not match the actual table schema and equality-delete files cannot be projected.
Solutions
- Compare each primary key column name against the Iceberg table schema (iceberg_table_metadata.current_schema()) and fix the sink's PRIMARY KEY to use exact existing column names
- Check for case sensitivity: Iceberg field lookup is exact, so normalize quoting/casing in the sink definition
- Re-create the sink against the correct table, or re-evolve the table to restore the dropped PK column
- Add a pre-flight validation in the sink creation path that resolves all PK columns to field ids before the writer is built
Example fix
// before
iceberg_schema.field_id_by_name(column_name).ok_or_else(|| {
SinkError::Config(anyhow!("Primary key column {} not found in Iceberg schema", column_name))
})
// after
iceberg_schema.field_id_by_name(column_name).ok_or_else(|| {
let known = iceberg_schema.fields().iter().map(|f| f.name()).collect::<Vec<_>>().join(", ");
SinkError::Config(anyhow!(
"Primary key column {} not found in Iceberg schema (available: [{}])",
column_name,
known
))
}) Defensive patterns
Strategy: validation
Validate before calling
for name in primary_key_column_names {
if iceberg_schema.field_id_by_name(name).is_none() {
return Err(format!("PK column {} missing from Iceberg schema", name));
}
} Type guard
fn pk_columns_resolve(names: &[String], schema: &iceberg::spec::Schema) -> bool {
names.iter().all(|n| schema.field_id_by_name(n).is_some())
} Try / catch
let ids = resolve_equality_delete_field_ids(&pk_names, &schema)
.map_err(|e| format!("sink PK does not match table schema: {e}"))?; Prevention
- Validate PK column names against the Iceberg schema at sink creation time
- Watch for case sensitivity when quoting column names in SQL
- Re-check the schema after any Iceberg table evolution (ALTER/DROP COLUMN)
- Keep the sink definition and the target table in the same migration/change review
When it happens
Trigger: Calling build_upsert for an Iceberg sink whose primary_key_column_names include a column missing (or renamed/case-mismatched) in table.metadata().current_schema(); sinking into a table evolved to drop a PK column.
Common situations: Typos or case differences between PRIMARY KEY in the CREATE SINK statement and the Iceberg table columns; pointing the sink at a different table than planned; Iceberg schema evolution removing a key column after sink creation.
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
- error from iceberg library
- no value find in sink schema, index is
- The schema of iceberg equality delete file must be…
- bounded compaction branch
- bounded compaction head sequence
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/0d107a8c307edb75.
Report an issue: GitHub.
Appendix: source
Thrown at src/connector/src/sink/iceberg/writer.rs:316
) -> Self {
Self::Created(IcebergSinkWriterArgs {
config,
sink_param,
writer_param,
upsert_primary_key_column_names,
})
}
}
pub(super) fn resolve_equality_delete_field_ids(
primary_key_column_names: &[String],
iceberg_schema: &iceberg::spec::Schema,
) -> Result<Vec<i32>> {
primary_key_column_names
.iter()
.map(|column_name| {
iceberg_schema.field_id_by_name(column_name).ok_or_else(|| {
SinkError::Config(anyhow!(
"Primary key column {} not found in Iceberg schema",
column_name
))
})
})
.collect()
}
impl IcebergSinkWriterInner {
pub fn build_append_only(
config: &IcebergConfig,
table: Table,
writer_param: &SinkWriterParam,
) -> Result<Self> {
let SinkWriterParam {
extra_partition_col_idx,
actor_id,
sink_id,View on GitHub (pinned to 6469eb736d)