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
- Inspect the upstream executor feeding the resolver side to see why non-insert ops are emitted.
- Verify executor wiring/version compatibility for the iceberg pk-index sink fragment.
- If replay/recovery produced this, restart the fragment from a clean checkpoint.
- 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
- Verify the resolver-side upstream is configured for insert-only output.
- Add a debug assertion on op type in executor tests.
- Keep executor topology changes covered by fragment integration tests.
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
- iceberg pk-index merger
- iceberg pk-index writer
- iceberg pk-index writer
- adlsgen2.authority_host does not parse as a URL
- adlsgen2.authority_host must not contain a path component
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)