risingwavelabs/risingwave · error

end of upstream

Error message

end of upstream

What it means

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.

Solutions

  1. Check meta service logs for upstream actor failure/failover around the time of the error and fix the root cause of the actor exit.
  2. Verify the backfill fragment's upstream dispatcher channels are correctly rebuilt after actor migration or scale-in.
  3. 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.

Example fix

// before
.ok_or_else(|| anyhow!("end of upstream"))?;
// after
// Optional observability improvement:
.ok_or_else(|| anyhow!("end of upstream (actor_id={}, fragment_id={})", self.actor_context.id, self.actor_context.fragment_id))?;
Defensive patterns

Strategy: try-catch

Try / catch

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()) }

Prevention

When it happens

Trigger: 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.

Common situations: 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.

Related errors


AI-assisted analysis of risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/d9081f9ad225624b. Report an issue: GitHub.

Appendix: source

Thrown at src/stream/src/executor/backfill/snapshot_backfill/executor.rs:813

                        warn!(pending_barrier = ?self.upstream_pending_barriers, "not polling upstream but timeout");
                        return pending().await;
                    }
                    self.consume_until_next_checkpoint_barrier().await?;
                } {
                    break e;
                }
            }
        }
    }

    /// Consume the upstream until seeing the next barrier.
    async fn consume_until_next_checkpoint_barrier(&mut self) -> StreamExecutorResult<()> {
        loop {
            let msg: DispatcherMessage = self
                .upstream
                .try_next()
                .await?
                .ok_or_else(|| anyhow!("end of upstream"))?;
            match msg {
                DispatcherMessage::Chunk(chunk) => {
                    self.is_polling_epoch_data = true;
                    self.consume_upstream_row_count
                        .inc_by(chunk.cardinality() as _);
                }
                DispatcherMessage::Barrier(barrier) => {
                    let is_checkpoint = barrier.kind.is_checkpoint();
                    self.upstream_pending_barriers.add(barrier);
                    if is_checkpoint {
                        self.is_polling_epoch_data = false;
                        break;
                    } else {
                        self.is_polling_epoch_data = true;
                    }
                }
                DispatcherMessage::Watermark(_) => {
                    self.is_polling_epoch_data = true;

View on GitHub (pinned to 6469eb736d)