risingwavelabs/risingwave · error · SinkError::Iceberg

error computing partition type

Error message

error computing partition type: {e}

What it means

resolve_partition_type converts a PartitionSpec into its partition StructType using the given schema; any iceberg-rust error from partition_type() is wrapped into this SinkError::Iceberg error. It typically means the spec's source field ids do not resolve against the provided schema.

Solutions

  1. Pass the schema that actually contains the partition source fields (e.g. the file's schema, not the newest table schema)
  2. Confirm the partition columns still exist in the sink's schema after any schema evolution
  3. Log/print the underlying iceberg error to identify which field id fails to resolve

Example fix

// before
let ptype = resolve_partition_type(&table, spec_id, &current_schema)?;
// after
let file_schema = table.metadata().schema_by_id(file.schema_id).unwrap();
let ptype = resolve_partition_type(&table, spec_id, file_schema)?;
Defensive patterns

Strategy: try-catch

Validate before calling

fn spec_fields_in_schema(spec_id: i32, schema: &Schema) -> bool {
    // verify every source field id of the spec resolves in schema
    spec.field_ids().iter().all(|id| schema.field_by_id(*id).is_some())
}

Try / catch

let ptype = resolve_partition_type(&table, spec_id, schema)
    .with_context(|| format!("computing partition type for spec {spec_id}"))?;

Prevention

When it happens

Trigger: Calling resolve_partition_type(table, spec_id, schema) where the partition spec references fields (source ids) missing from `schema`, or the iceberg-rust partition_type computation fails on an incompatible field type (e.g. spec written against an older schema with since-changed types).

Common situations: Schema evolution (dropped/renamed partition columns) between when data files were written and committed; passing the wrong schema variant (e.g. current schema vs. file's schema) to the resolver; version upgrades in iceberg-rust changing validation strictness.

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

Appendix: source

Thrown at src/connector/src/sink/iceberg/writer.rs:1030

/// error message. Shared spec-lookup mechanic for the pk-index merger, commit
/// coordinator, and sink commit paths.
pub fn resolve_partition_spec(table: &Table, spec_id: i32) -> Result<PartitionSpecRef> {
    table
        .metadata()
        .partition_spec_by_id(spec_id)
        .cloned()
        .ok_or_else(|| SinkError::Iceberg(anyhow!("partition spec {} not found", spec_id)))
}

/// Resolves the partition [`StructType`] for the given `spec_id` against `schema`.
///
/// `schema` is passed explicitly (rather than read from the table) so callers can
/// preserve their chosen schema, and `spec_id` is passed explicitly so callers can
/// preserve their chosen spec (e.g. a file's own spec vs. the default spec).
pub fn resolve_partition_type(table: &Table, spec_id: i32, schema: &Schema) -> Result<StructType> {
    resolve_partition_spec(table, spec_id)?
        .partition_type(schema)
        .map_err(|e| SinkError::Iceberg(anyhow!(e)))
}

/// Truncate large column statistics from `DataFile` BEFORE serialization.
///
/// This function directly modifies `DataFile`'s `lower_bounds` and `upper_bounds`
/// to remove entries that exceed `MAX_COLUMN_STAT_SIZE`.
///
/// # Arguments
/// * `data_file` - A `DataFile` to process
///
/// # Returns
/// The modified `DataFile` with large statistics truncated
pub fn truncate_datafile(mut data_file: DataFile) -> DataFile {
    // Process lower_bounds - remove entries with large values
    data_file.lower_bounds_mut().retain(|field_id, datum| {
        // Use to_bytes() to get the actual binary size without JSON serialization overhead
        let size = match datum.to_bytes() {
            Ok(bytes) => bytes.len(),

View on GitHub (pinned to 6469eb736d)