risingwavelabs/risingwave · error · StreamExecutorError
iceberg pk-index writer
Error message
iceberg pk-index writer {} resolver input closed before switch-to-input for task {} What it means
The writer was resolving an Iceberg compaction task and waiting on the remap input for its Phase::End barrier, but the remap input stream ended (channel closed) before delivering it. Without that barrier the writer cannot transition to the replacement input, so it fails with this error.
Solutions
- Check meta/log files for the upstream remap actor's status around the failure time — restart or failure of that actor is the usual cause.
- Retry the compaction task or the sink after the cluster is healthy; the task state is tracked so it can be re-driven.
- If this recurs during admin operations (scaling, sink recreation), avoid mutating the sink/fragment while a compaction task is in flight.
- Capture the sink_id and task_id from the message and file an issue with the actor logs if no upstream failure is visible.
Defensive patterns
Strategy: retry
Validate before calling
// Check the remap upstream actor health before/at compaction start
async fn remap_input_healthy(meta: &MetaClient, fragment_id: FragmentId) -> bool {
meta.actor_status(fragment_id).await.map(|s| s.is_running()).unwrap_or(false)
} Try / catch
// Wrap compaction-driving operations and retry after upstream recovery
match run_compaction(task_id).await {
Err(e) if e.to_string().contains("resolver input closed") => {
wait_for_cluster_healthy().await;
retry_compaction(task_id).await?;
}
other => other?,
} Prevention
- Avoid administrative operations (scaling, sink recreation, config restarts) while an iceberg compaction task is in flight.
- Monitor actor liveness so upstream failures are caught and recovered before the writer observes a closed channel.
- Retry failed compaction tasks once the cluster is healthy; the task state is durable.
When it happens
Trigger: In execute_resolving_right, the #[for_await] loop over resolver_input completes normally (stream exhausted, no error) before any Message::Barrier matching the Phase::End compaction barrier is seen.
Common situations: The upstream remap actor crashed or was scaled/cancelled mid-compaction (e.g., during a config change, failover, or `ALTER SINK` restart); a network or channel teardown between the remap operator and the writer; cluster shutdown racing an in-flight compaction task.
Related errors
- iceberg pk-index writer
- iceberg pk-index writer
- barrier reader closed unexpectedly
- build parquet stream for
- cannot get truncated epoch
AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11).
Data as JSON: /api/errors/c8a15ca2f33c8b88.
Report an issue: GitHub.
Appendix: source
Thrown at src/stream/src/executor/iceberg_with_pk_index/writer.rs:515
Message::Watermark(_) => bail!(
"iceberg pk-index writer {} received watermark on remap input while resolving task {}",
self.sink_id,
task_id
),
Message::Barrier(barrier) => {
self.validate_compaction_barrier(
&barrier,
task_id,
Phase::End,
begin_epoch.curr,
)?;
self.mode = WriterInputMode::AligningReplacementInput { task_id, barrier };
return Ok(());
}
}
}
bail!(
"iceberg pk-index writer {} resolver input closed before switch-to-input for task {}",
self.sink_id,
task_id
);
}
#[try_stream(ok = Message, error = StreamExecutorError)]
async fn execute_aligning_replacement_input<'a>(
&'a mut self,
input: &'a mut BoxedMessageStream,
task_id: IcebergCompactionTaskId,
expected: Barrier,
) {
#[for_await]
for msg in input {
match msg? {
Message::Chunk(_) => bail!(
"iceberg pk-index writer {} received a chunk from replacement input before its initial barrier for task {}",View on GitHub (pinned to 6469eb736d)