risingwavelabs/risingwave · error

iceberg pk-index sink reports disagree on schema_id

Error message

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

What it means

All worker reports in one Iceberg pk-index sink commit must reference the same table schema. align_report_id records the first observed schema_id and bails if a later report disagrees, because committing files written under different schemas would corrupt the Iceberg table. The error reports both conflicting ids.

Solutions

  1. Identify the lagging worker from logs and restart it so it reloads the current table schema.
  2. Ensure schema-change notifications fully propagate; wait for the next epoch so all workers align on the new schema.
  3. Check the Iceberg table's schema history to confirm the expected schema_id and reconcile workers pinned to an outdated snapshot.
  4. If a worker is stuck with a cached schema, restart the sink fragment.
Defensive patterns

Strategy: validation

Validate before calling

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

Try / catch

// retry once after workers refresh schema
if let Err(e) = commit().await {
    if e.to_string().contains("disagree on schema_id") {
        refresh_worker_schemas().await;
        commit().await?;
    } else {
        return Err(e);
    }
}

Prevention

When it happens

Trigger: pre_commit_epoch -> aggregate_reports -> align_report_id when two reports for the same sink/epoch carry different schema_id values, e.g. one worker still on the pre-evolution schema and another on the post-evolution one.

Common situations: Schema evolution landing mid-epoch while some writers/mergers still run with the old cached schema; stale workers that missed the schema-change notification; replayed old reports merged with fresh ones.

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/c4c3ff3e4ba8461d. Report an issue: GitHub.

Appendix: source

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

    Ok(IcebergPkIndexSinkAggResult {
        schema_id: shared_schema_id.unwrap(),
        partition_spec_id: shared_partition_spec_id.unwrap(),
        data_files,
        delete_files,
        overwrite_files,
    })
}

fn align_report_id(
    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),
        _ => {}

View on GitHub (pinned to 6469eb736d)