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

  1. 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
  2. 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
  3. Read the wrapped iceberg library error to see which specific column id failed validation
  4. 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

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


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)