{"record":{"id":"c8a15ca2f33c8b88","repo":"risingwavelabs/risingwave","slug":"iceberg-pk-index-writer-resolver-input-closed-b","errorCode":null,"errorMessage":"iceberg pk-index writer {} resolver input closed before switch-to-input for task {}","messagePattern":"iceberg pk-index writer (.+?) resolver input closed before switch-to-input for task (.+?)","errorType":"exception","errorClass":"StreamExecutorError","httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/iceberg_with_pk_index/writer.rs","lineNumber":515,"sourceCode":"                Message::Watermark(_) => bail!(\n                    \"iceberg pk-index writer {} received watermark on remap input while resolving task {}\",\n                    self.sink_id,\n                    task_id\n                ),\n                Message::Barrier(barrier) => {\n                    self.validate_compaction_barrier(\n                        &barrier,\n                        task_id,\n                        Phase::End,\n                        begin_epoch.curr,\n                    )?;\n                    self.mode = WriterInputMode::AligningReplacementInput { task_id, barrier };\n                    return Ok(());\n                }\n            }\n        }\n\n        bail!(\n            \"iceberg pk-index writer {} resolver input closed before switch-to-input for task {}\",\n            self.sink_id,\n            task_id\n        );\n    }\n\n    #[try_stream(ok = Message, error = StreamExecutorError)]\n    async fn execute_aligning_replacement_input<'a>(\n        &'a mut self,\n        input: &'a mut BoxedMessageStream,\n        task_id: IcebergCompactionTaskId,\n        expected: Barrier,\n    ) {\n        #[for_await]\n        for msg in input {\n            match msg? {\n                Message::Chunk(_) => bail!(\n                    \"iceberg pk-index writer {} received a chunk from replacement input before its initial barrier for task {}\",","sourceCodeStart":497,"sourceCodeEnd":533,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/iceberg_with_pk_index/writer.rs#L497-L533","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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."],"exampleFix":null,"handlingStrategy":"retry","validationCode":"// Check the remap upstream actor health before/at compaction start\nasync fn remap_input_healthy(meta: &MetaClient, fragment_id: FragmentId) -> bool {\n    meta.actor_status(fragment_id).await.map(|s| s.is_running()).unwrap_or(false)\n}","typeGuard":null,"tryCatchPattern":"// Wrap compaction-driving operations and retry after upstream recovery\nmatch run_compaction(task_id).await {\n    Err(e) if e.to_string().contains(\"resolver input closed\") => {\n        wait_for_cluster_healthy().await;\n        retry_compaction(task_id).await?;\n    }\n    other => other?,\n}","preventionTips":["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."],"tags":["streaming","channel-closed","iceberg-sink","compaction"],"backgroundTag":"broken-pipe","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}