risingwavelabs/risingwave · error · StreamExecutorError

iceberg pk-index writer

Error message

iceberg pk-index writer {} expected resolver inserts, got {op:?}

What it means

The resolver side of the pk-index writer only expects Insert operations: the equi-join resolver produces the rows that must be inserted into the pk-index state table. Any Update/Delete op reaching apply_resolver_chunk violates the protocol and aborts with this error. The state table is append-only from this path, so non-inserts are unrecoverable here.

Solutions

  1. Inspect the upstream executor feeding the resolver side to see why non-insert ops are emitted.
  2. Verify executor wiring/version compatibility for the iceberg pk-index sink fragment.
  3. If replay/recovery produced this, restart the fragment from a clean checkpoint.
  4. File a bug including the op variant and actor id if it occurs with standard pipelines.

Example fix

// before
if op != Op::Insert {
    bail!("... expected resolver inserts, got {op:?}");
}
// after
// keep as-is: fix the upstream instead — e.g. ensure the materialize/resolver input
// is configured with INSERT-only output (no update/delete propagation).
Defensive patterns

Strategy: validation

Validate before calling

// filter or assert at the producer side feeding the resolver input
for (op, row) in chunk.rows() {
    assert_eq!(op, Op::Insert, "resolver input must be insert-only");
}

Type guard

fn is_insert_only(chunk: &StreamChunk) -> bool {
    chunk.rows().all(|(op, _)| op == Op::Insert)
}

Try / catch

if let Err(e) = writer.execute_resolving_right(chunk).await {
    if e.to_string().contains("expected resolver inserts") {
        error!("non-insert op reached resolver input; check upstream executor wiring");
    }
    return Err(e);
}

Prevention

When it happens

Trigger: Raised in apply_resolver_chunk (called from execute_resolving_right) when a StreamChunk row from the resolving (right) side carries Op::UpdateDelete, Op::UpdateInsert, or Op::Delete instead of Op::Insert.

Common situations: Upstream executor changes emitting CDC-style updates into the resolver side; incorrect join/executor wiring that passes mutating ops to the writer's resolver input; data replay path differences after recovery.

Understand the failure class

Background: Invalid enum value errors: "Unknown type", "Invalid scope", "must be one of" — when a string is not on the library's allowed list — this error's family across 23 libraries.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/44079962ce2fdb5e. Report an issue: GitHub.

Appendix: source

Thrown at src/stream/src/executor/iceberg_with_pk_index/writer.rs:264

            for (pk, pos) in insert_pks.into_iter().zip_eq_fast(positions) {
                let mut index_row_data = Vec::with_capacity(pk_indices.len() + 2);
                for datum in pk.iter() {
                    index_row_data.push(datum);
                }
                index_row_data.push(Some(ScalarRefImpl::Utf8(&pos.path)));
                index_row_data.push(Some(ScalarRefImpl::Int64(pos.pos)));
                self.pk_index_state_table.insert(index_row_data.as_slice());
            }
        }

        self.delete_position_buffer = Some(delete_position_buffer);
        self.pk_index_state_table.try_flush().await?;
    }

    async fn apply_resolver_chunk(&mut self, chunk: StreamChunk) -> StreamExecutorResult<()> {
        for (op, row) in chunk.rows() {
            if op != Op::Insert {
                bail!(
                    "iceberg pk-index writer {} expected resolver inserts, got {op:?}",
                    self.sink_id
                );
            }
            self.pk_index_state_table.insert(row);
        }
        self.pk_index_state_table.try_flush().await?;
        Ok(())
    }

    #[try_stream(ok = Message, error = StreamExecutorError)]
    async fn checkpoint_barrier(&mut self, barrier: Barrier) {
        barrier.assume_no_update_vnode_bitmap(self.ctx.id)?;
        let mut metadata = None;
        if barrier.is_checkpoint() {
            if let Some(chunk) = self
                .delete_position_buffer
                .take()

View on GitHub (pinned to 6469eb736d)