risingwavelabs/risingwave · error

iceberg pk-index sink reports disagree on partition_spec_id

Error message

iceberg pk-index sink reports disagree on partition_spec_id: {} vs {}

What it means

Analogous to the schema check: every report in a single commit must agree on the table's partition_spec_id. align_report_id compares each report's id against the shared value and bails with both ids when they differ, preventing a commit that mixes files partitioned under different specs.

Solutions

  1. Identify the worker with the stale partition_spec_id from logs and restart it to reload current table metadata.
  2. Force a refresh of Iceberg table metadata on all workers so partition specs align before the next commit.
  3. Confirm the intended spec id on the Iceberg table and roll back an accidental spec change if unintended.
  4. Restart the sink fragment so all reports regenerate against the same spec.
Defensive patterns

Strategy: validation

Validate before calling

// before triggering commit, ensure all workers report the same partition spec id
let ids: HashSet<i32> = reports.iter().map(|r| r.partition_spec_id).collect();
if ids.len() > 1 {
    return Err(anyhow!("workers disagree on partition_spec_id: {:?}", ids));
}

Try / catch

// retry once after refreshing table metadata
if let Err(e) = commit().await {
    if e.to_string().contains("disagree on partition_spec_id") {
        refresh_table_metadata().await;
        commit().await?;
    } else {
        return Err(e);
    }
}

Prevention

When it happens

Trigger: pre_commit_epoch -> aggregate_reports -> align_report_id when two worker reports for the same sink/epoch carry different partition_spec_id values, typically around a partition-spec evolution.

Common situations: Partition spec evolution applied while workers still hold the old spec; a worker restored from an old checkpoint/metadata snapshot; replayed stale reports from before the spec change.

Understand the failure class

Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.

Related errors


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

Appendix: source

Thrown at src/meta/src/manager/iceberg_pk_index_sink/coordinator.rs:639

    schema_id: i32,
    partition_spec_id: i32,
    shared_schema_id: &mut Option<i32>,
    shared_partition_spec_id: &mut Option<i32>,
) -> Result<()> {
    match shared_schema_id {
        Some(prev) if *prev != schema_id => {
            bail!(
                "iceberg pk-index sink reports disagree on schema_id: {} vs {}",
                prev,
                schema_id
            );
        }
        None => *shared_schema_id = Some(schema_id),
        _ => {}
    }
    match shared_partition_spec_id {
        Some(prev) if *prev != partition_spec_id => {
            bail!(
                "iceberg pk-index sink reports disagree on partition_spec_id: {} vs {}",
                prev,
                partition_spec_id
            );
        }
        None => *shared_partition_spec_id = Some(partition_spec_id),
        _ => {}
    }
    Ok(())
}

fn encode_pre_commit_state(
    agg_result: &IcebergPkIndexSinkAggResult,
    snapshot_id: i64,
) -> Result<Vec<u8>> {
    let agg_result = serde_json::to_vec(agg_result)?;
    Ok(PbIcebergPkIndexPreCommitState {
        agg_result,

View on GitHub (pinned to 6469eb736d)