{"record":{"id":"d9081f9ad225624b","repo":"risingwavelabs/risingwave","slug":"end-of-upstream","errorCode":null,"errorMessage":"end of upstream","messagePattern":"end of upstream","errorType":"exception","errorClass":null,"httpStatus":null,"severity":"error","filePath":"src/stream/src/executor/backfill/snapshot_backfill/executor.rs","lineNumber":813,"sourceCode":"                        warn!(pending_barrier = ?self.upstream_pending_barriers, \"not polling upstream but timeout\");\n                        return pending().await;\n                    }\n                    self.consume_until_next_checkpoint_barrier().await?;\n                } {\n                    break e;\n                }\n            }\n        }\n    }\n\n    /// Consume the upstream until seeing the next barrier.\n    async fn consume_until_next_checkpoint_barrier(&mut self) -> StreamExecutorResult<()> {\n        loop {\n            let msg: DispatcherMessage = self\n                .upstream\n                .try_next()\n                .await?\n                .ok_or_else(|| anyhow!(\"end of upstream\"))?;\n            match msg {\n                DispatcherMessage::Chunk(chunk) => {\n                    self.is_polling_epoch_data = true;\n                    self.consume_upstream_row_count\n                        .inc_by(chunk.cardinality() as _);\n                }\n                DispatcherMessage::Barrier(barrier) => {\n                    let is_checkpoint = barrier.kind.is_checkpoint();\n                    self.upstream_pending_barriers.add(barrier);\n                    if is_checkpoint {\n                        self.is_polling_epoch_data = false;\n                        break;\n                    } else {\n                        self.is_polling_epoch_data = true;\n                    }\n                }\n                DispatcherMessage::Watermark(_) => {\n                    self.is_polling_epoch_data = true;","sourceCodeStart":795,"sourceCodeEnd":831,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/stream/src/executor/backfill/snapshot_backfill/executor.rs#L795-L831","documentation":"The upstream channel of the snapshot backfill executor closed, so `try_next()` returned None while waiting for the next checkpoint barrier or data chunk. The executor treats a closed upstream as unrecoverable because streaming executors are expected to have a live upstream that only ends via barrier-driven termination.","triggerScenarios":"All upstream `DispatcherMessage` senders are dropped while `consume_until_next_checkpoint_barrier` is polling — e.g. the upstream actor fails or is cancelled mid-backfill, or the dispatcher channel is closed during migration/scale-in.","commonSituations":"Upstream actor crash or failover during snapshot backfill; actor rescheduling on a different node without proper channel rebuild; cluster shutdown while a backfill is still consuming.","solutions":["Check meta service logs for upstream actor failure/failover around the time of the error and fix the root cause of the actor exit.","Verify the backfill fragment's upstream dispatcher channels are correctly rebuilt after actor migration or scale-in.","Retry the MV/table creation; if it reproduces, capture the full trace and file an issue — a closed upstream mid-backfill usually indicates a meta/actor lifecycle bug."],"exampleFix":"// before\n.ok_or_else(|| anyhow!(\"end of upstream\"))?;\n// after\n// Optional observability improvement:\n.ok_or_else(|| anyhow!(\"end of upstream (actor_id={}, fragment_id={})\", self.actor_context.id, self.actor_context.fragment_id))?;","handlingStrategy":"try-catch","validationCode":null,"typeGuard":null,"tryCatchPattern":"match self.upstream.try_next().await { Ok(Some(msg)) => ..., Ok(None) => { warn!(\"upstream closed\"); return Err(anyhow!(\"end of upstream\").into()); }, Err(e) => return Err(e.into()) }","preventionTips":["Monitor actor liveness and restart failed upstream actors promptly","Avoid cancelling actors mid-backfill during maintenance","Alert on dispatcher channel closes during scale-in/migration"],"tags":["streaming","backfill","channel-closed"],"backgroundTag":"unexpected-response-shape","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"}