risingwavelabs/risingwave · error
missing primary-key columns in compaction resolver
Error message
missing primary-key columns in compaction resolver
What it means
The Iceberg compaction resolver conversion derives primary key indices from `node.pk_columns` and requires at least one entry. An empty list means the resolver cannot identify row identity for resolving compaction results, so `new_boxed_executor` fails early with this message.
Solutions
- Define primary key columns on the Iceberg sink/table so pk_columns are populated in the proto.
- Upgrade frontend/meta so pk_columns are serialized for iceberg_with_pk_index nodes.
- Recreate the sink materialization to regenerate the fragment with PK columns.
Example fix
-- before CREATE SINK s FROM mv; -- no primary key -- after CREATE SINK s FROM mv WITH PRIMARY KEY (id); -- pk_columns now non-empty
Defensive patterns
Strategy: validation
Validate before calling
// before building the node proto assert!(!pk_columns.is_empty(), "iceberg with-pk-index sink requires primary key columns");
Type guard
fn has_pk(node: &StreamNode) -> bool { !node.pk_columns.is_empty() } Prevention
- Always define a primary key for iceberg_with_pk_index sinks
- Add plan validation that pk_columns is non-empty at planning time
- Check proto serialization after frontend changes
When it happens
Trigger: Building the compaction resolver executor when `node.pk_columns` is empty (or all entries map to zero data-file indices yielding an empty vec), i.e. the proto node was produced without PK columns for a with-pk-index Iceberg sink.
Common situations: Planning an Iceberg sink with PK index but omitting primary key definition; frontend serialization bug dropping pk_columns; older plan versions lacking the pk_columns field.
Understand the failure class
Background: "missing required argument" and "the following required arguments were not provided": what required-argument errors mean and how to fix them — this error's family across 20 libraries.
Related errors
- compaction resolver PK column missing column_desc
- error from iceberg library
- failed to receive the first barrier, actor_id
- Iceberg metadata relations are not supported in streaming…
- iceberg pk-index writer
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/1e20de0b7a399c0d.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/from_proto/iceberg_with_pk_index/compaction_resolver.rs:57
assert!(
params.input.is_empty(),
"compaction resolver executor should not have input"
);
let sink_id = node.sink_id;
let properties_with_secret = LocalSecretManager::global()
.fill_secrets(node.properties.clone(), node.secret_refs.clone())?;
let iceberg_config = IcebergConfig::from_btreemap(properties_with_secret)
.map_err(|err| StreamExecutorError::from((err, sink_id)))?;
let pk_indices = node
.pk_columns
.iter()
.map(|column| column.data_file_index as usize)
.collect::<Vec<_>>();
if pk_indices.is_empty() {
return Err(anyhow!("missing primary-key columns in compaction resolver").into());
}
let pk_data_types = node
.pk_columns
.iter()
.map(|column| {
column
.column_desc
.as_ref()
.map(ColumnDesc::from)
.map(|column| column.data_type)
.ok_or_else(|| anyhow!("compaction resolver PK column missing column_desc"))
})
.collect::<Result<Vec<_>, _>>()?;
let barrier_receiver = params
.local_barrier_manager
.subscribe_barrier(params.actor_context.id);View on GitHub (pinned to 6469eb736d)