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
- Pass the schema that actually contains the partition source fields (e.g. the file's schema, not the newest table schema)
- Confirm the partition columns still exist in the sink's schema after any schema evolution
- 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, ¤t_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
- Pass the schema version that contains the spec's source fields
- Track schema/spec evolution: on change, verify sink state still resolves
- Log the inner iceberg error to identify the failing field id
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
- auto schema refresh sink must have only one fragment, but…
- Can't create iceberg sink write result from empty data!
- Can't find data
- Can't find schema by id
- Cannot find
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)