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

  1. Define primary key columns on the Iceberg sink/table so pk_columns are populated in the proto.
  2. Upgrade frontend/meta so pk_columns are serialized for iceberg_with_pk_index nodes.
  3. 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

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


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)