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
- Inspect the underlying conversion error in the log to see which field/partition value failed validation.
- 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.
- 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
- Avoid external ALTER TABLE / partition-spec changes on sink-managed tables.
- Keep pending epochs short so schema skew windows stay small.
- Add a pre-commit check that table schema_id/partition_spec_id match the reports.
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
- apply iceberg pk-index sink overwrite_files action
- backfill iceberg pk-index sink delete files failed…
- error converting StreamChunk to Arrow RecordBatch
- Expect return type but got , RisingWave return type is …
- Failed to update iceberg table.
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_filesView on GitHub (pinned to 6469eb736d)