risingwavelabs/risingwave · error

materialize pk-index sink SerializedDataFile

Error message

materialize pk-index sink SerializedDataFile

What it means

During commit, serialized data files persisted in the aggregation result must be materialized back into Iceberg `DataFile`s against the table's current schema and partition type. `SerializedDataFile::try_into` validates partition values and schema fields; a mismatch throws this error, wrapped as a `CommitError::Commit`.

Solutions

  1. Inspect the underlying conversion error in the log to see which field/partition value failed validation.
  2. Check whether the Iceberg table's schema or partition spec changed externally after the report was produced; align/revert the external change or recompute the commit.
  3. Abort the stale pending epoch so it can be re-run with reports matching the current table schema.

Example fix

// avoid external evolution of a sink-managed table:
// before: ALTER TABLE db.tbl SET PARTITION SPEC (...); -- by external engine
// after: coordinate spec changes with RisingWave sink downtime
Defensive patterns

Strategy: try-catch

Validate before calling

// verify schema/spec stability before commit
anyhow::ensure!(
    table.metadata().current_schema().schema_id() == merged.schema_id,
    "table schema changed since pre-commit"
);
anyhow::ensure!(
    table.metadata().default_partition_spec_id() == merged.partition_spec_id,
    "partition spec changed since pre-commit"
);

Try / catch

match serialized.clone().try_into(spec_id, &partition_type, schema) {
    Ok(df) => { /* use df */ }
    Err(e) => return Err(CommitError::Commit(anyhow!(e).context("materialize pk-index sink SerializedDataFile"))),
}

Prevention

When it happens

Trigger: In `commit_one_epoch`, the persisted `SerializedDataFile` cannot convert with the current `partition_spec_id` / `partition_type` / schema — e.g. table schema evolved or partition spec changed between pre-commit persistence and commit, or the blob was written against a different table version.

Common situations: External schema evolution (ALTER TABLE on the Iceberg table by another engine) while a commit is pending; recovered pending rows from an older table snapshot; partition spec replaced in the catalog.

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

Appendix: source

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

                    .snapshots()
                    .any(|s| s.snapshot_id() == snapshot_id)
                {
                    return Ok(table);
                }

                let schema = table.metadata().current_schema();
                let partition_type =
                    resolve_partition_type(&table, merged.partition_spec_id, schema)
                        .map_err(|e| CommitError::Commit(anyhow!(e)))?;

                let materialize =
                    |serialized: &SerializedDataFile| -> Result<DataFile, CommitError> {
                        serialized
                            .clone()
                            .try_into(merged.partition_spec_id, &partition_type, schema)
                            .map_err(|err| {
                                CommitError::Commit(
                                    anyhow!(err)
                                        .context("materialize pk-index sink SerializedDataFile"),
                                )
                            })
                    };

                // Reuse the files already materialized during pre-commit backfill when available;
                // otherwise (recovery / unpartitioned) materialize from the persisted form once.
                let add_files: Vec<DataFile> = match materialized_add_files {
                    Some(add_files) => add_files,
                    None => merged
                        .data_files
                        .iter()
                        .chain(merged.delete_files.iter())
                        .map(&materialize)
                        .collect::<Result<Vec<_>, _>>()?,
                };
                let overwrite_files: Vec<DataFile> = merged
                    .overwrite_files

View on GitHub (pinned to 6469eb736d)