risingwavelabs/risingwave · error · SinkError::Iceberg
error from iceberg library
Error message
error from iceberg library: {err} (EqualityDeleteWriterConfig::new failed) What it means
EqualityDeleteWriterConfig::new(unique_column_ids, current_schema) failed in iceberg-rs while setting up the equality-delete writer. This constructor validates that the given column ids exist in the schema and are usable as an equality-delete key set; any mismatch aborts build_upsert with SinkError::Iceberg.
Solutions
- Re-derive unique_column_ids from the current table schema via iceberg_schema.field_id_by_name for each PK column and confirm none are None
- Verify the sink's PRIMARY KEY columns still exist in the Iceberg table (schema evolution may have removed them) and update the sink or table accordingly
- Read the wrapped iceberg library error to see which specific column id failed validation
- Ensure unique_column_ids is non-empty before constructing EqualityDeleteWriterConfig
Example fix
// before
let eq_del_config = EqualityDeleteWriterConfig::new(
unique_column_ids.clone(),
table.metadata().current_schema().clone(),
)
.map_err(|err| SinkError::Iceberg(anyhow!(err)))?;
// after
anyhow::ensure!(!unique_column_ids.is_empty(), "equality delete needs >= 1 column id");
let eq_del_config = EqualityDeleteWriterConfig::new(
unique_column_ids.clone(),
table.metadata().current_schema().clone(),
)
.map_err(|err| {
SinkError::Iceberg(anyhow::anyhow!(
"EqualityDeleteWriterConfig::new failed: {err}; ids = {:?}",
unique_column_ids
))
})? Defensive patterns
Strategy: validation
Validate before calling
anyhow::ensure!(!unique_column_ids.is_empty(), "equality delete config requires at least one column id");
let schema_field_ids: std::collections::HashSet<i32> =
table.metadata().current_schema().fields().map(|f| f.id).collect();
anyhow::ensure!(
unique_column_ids.iter().all(|id| schema_field_ids.contains(id)),
"equality delete ids not present in schema"
); Try / catch
let eq_del_config = EqualityDeleteWriterConfig::new(
unique_column_ids.clone(),
table.metadata().current_schema().clone(),
)
.map_err(|err| SinkError::Iceberg(anyhow!("EqualityDeleteWriterConfig::new failed: {err} (ids: {:?})", unique_column_ids)))?; Prevention
- Re-resolve PK column ids from the current table schema at writer build time, never cache across schema evolutions
- Confirm all PK columns still exist after any Iceberg ALTER/DROP COLUMN
- Ensure the PK set is non-empty before building the equality-delete writer
- Include the offending column ids in the wrapped error message
When it happens
Trigger: Calling build_upsert when unique_column_ids (derived from the sink's primary key via resolve_equality_delete_field_ids) contains an id not present in table.metadata().current_schema(), an empty id set, or a projected schema inconsistent with the ids.
Common situations: Iceberg schema evolution (ids dropped/changed) after the sink computed its column ids; stale/mismatched sink configuration pointing at a table whose key columns differ; empty PK mapping producing an empty id list rejected by the library.
Understand the failure class
Background: "Invalid value" and "allowed values are" config errors: what your library rejected and how to fix it — this error's family across 41 libraries.
Related errors
- Primary key column not found in Iceberg schema
- The schema of iceberg equality delete file must be…
- bounded compaction branch
- bounded compaction head sequence
- bounded compaction is not supported for copy-on-write tasks
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/494a109b2a8f0b82.
Report an issue: GitHub.
Appendix: source
Thrown at src/connector/src/sink/iceberg/writer.rs:578
table.file_io().clone(),
DefaultLocationGenerator::new(table.metadata())
.map_err(|err| SinkError::Iceberg(anyhow!(err)))?,
DefaultFileNameGenerator::new(
writer_param.actor_id.to_string(),
Some(format!("pos-del-{}", unique_uuid_suffix)),
iceberg::spec::DataFileFormat::Parquet,
),
);
PositionDeleteWriterBuilderType::PositionDelete(PositionDeleteFileWriterBuilder::new(
rolling_writer_builder,
))
};
let equality_delete_builder = {
let eq_del_config = EqualityDeleteWriterConfig::new(
unique_column_ids.clone(),
table.metadata().current_schema().clone(),
)
.map_err(|err| SinkError::Iceberg(anyhow!(err)))?;
let parquet_writer_builder = ParquetWriterBuilder::new(
parquet_writer_properties,
Arc::new(
arrow_schema_to_schema(eq_del_config.projected_arrow_schema_ref())
.map_err(|err| SinkError::Iceberg(anyhow!(err)))?,
),
);
let rolling_writer_builder = RollingFileWriterBuilder::new(
parquet_writer_builder,
(config.target_file_size_mb() * 1024 * 1024) as usize,
table.file_io().clone(),
DefaultLocationGenerator::new(table.metadata())
.map_err(|err| SinkError::Iceberg(anyhow!(err)))?,
DefaultFileNameGenerator::new(
writer_param.actor_id.to_string(),
Some(format!("eq-del-{}", unique_uuid_suffix)),
iceberg::spec::DataFileFormat::Parquet,
),View on GitHub (pinned to 6469eb736d)