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

  1. 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.
  2. Retry the compaction task or the sink after the cluster is healthy; the task state is tracked so it can be re-driven.
  3. If this recurs during admin operations (scaling, sink recreation), avoid mutating the sink/fragment while a compaction task is in flight.
  4. 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

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


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)