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
- 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.
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
- 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
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
- barrier reader closed unexpectedly
- cannot get truncated epoch
- end of barrier receiver
- end of new output request
- end of stream
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)